feat(ui): finalize E-TIB modernization with footer redesign, video optimization, and TS fixes

Former-commit-id: 67ac02c8404cc66893fdf97308574701cca6000c
This commit is contained in:
2026-05-07 10:43:20 +02:00
parent f2fdea96ec
commit 5439df72ea
3712 changed files with 932332 additions and 1867 deletions
@@ -0,0 +1 @@
../../real-require@0.2.0/node_modules/real-require
@@ -0,0 +1,21 @@
MIT License
Copyright (c) 2021 Matteo Collina
Permission is hereby granted, free of charge, to any person obtaining a copy
of this software and associated documentation files (the "Software"), to deal
in the Software without restriction, including without limitation the rights
to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
copies of the Software, and to permit persons to whom the Software is
furnished to do so, subject to the following conditions:
The above copyright notice and this permission notice shall be included in all
copies or substantial portions of the Software.
THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE
SOFTWARE.
@@ -0,0 +1,85 @@
'use strict'
const bench = require('fastbench')
const SonicBoom = require('sonic-boom')
const ThreadStream = require('.')
const Console = require('console').Console
const fs = require('fs')
const { join } = require('path')
const core = fs.createWriteStream('/dev/null')
const fd = fs.openSync('/dev/null', 'w')
const sonic = new SonicBoom({ fd })
const sonicSync = new SonicBoom({ fd, sync: true })
const out = fs.createWriteStream('/dev/null')
const dummyConsole = new Console(out)
const threadStreamSync = new ThreadStream({
filename: join(__dirname, 'test', 'to-file.js'),
workerData: { dest: '/dev/null' },
bufferSize: 4 * 1024 * 1024,
sync: true
})
const threadStreamAsync = new ThreadStream({
filename: join(__dirname, 'test', 'to-file.js'),
workerData: { dest: '/dev/null' },
bufferSize: 4 * 1024 * 1024,
sync: false
})
const MAX = 10000
let str = ''
for (let i = 0; i < 100; i++) {
str += 'hello'
}
setTimeout(doBench, 100)
const run = bench([
function benchThreadStreamSync (cb) {
for (let i = 0; i < MAX; i++) {
threadStreamSync.write(str)
}
setImmediate(cb)
},
function benchThreadStreamAsync (cb) {
threadStreamAsync.once('drain', cb)
for (let i = 0; i < MAX; i++) {
threadStreamAsync.write(str)
}
},
function benchSonic (cb) {
sonic.once('drain', cb)
for (let i = 0; i < MAX; i++) {
sonic.write(str)
}
},
function benchSonicSync (cb) {
sonicSync.once('drain', cb)
for (let i = 0; i < MAX; i++) {
sonicSync.write(str)
}
},
function benchCore (cb) {
core.once('drain', cb)
for (let i = 0; i < MAX; i++) {
core.write(str)
}
},
function benchConsole (cb) {
for (let i = 0; i < MAX; i++) {
dummyConsole.log(str)
}
setImmediate(cb)
}
], 1000)
function doBench () {
run(function () {
run(function () {
// TODO figure out why it does not shut down
process.exit(0)
})
})
}
@@ -0,0 +1,10 @@
'use strict'
const neostandard = require('neostandard')
module.exports = neostandard({
ignores: [
'test/ts/**/*',
'test/syntax-error.mjs'
]
})
@@ -0,0 +1,541 @@
'use strict'
const { version } = require('./package.json')
const { EventEmitter } = require('events')
const { Worker } = require('worker_threads')
const { join } = require('path')
const { pathToFileURL } = require('url')
const { wait } = require('./lib/wait')
const {
WRITE_INDEX,
READ_INDEX
} = require('./lib/indexes')
const buffer = require('buffer')
const assert = require('assert')
const kImpl = Symbol('kImpl')
// V8 limit for string size
const MAX_STRING = buffer.constants.MAX_STRING_LENGTH
class FakeWeakRef {
constructor (value) {
this._value = value
}
deref () {
return this._value
}
}
class FakeFinalizationRegistry {
register () {}
unregister () {}
}
// Currently using FinalizationRegistry with code coverage breaks the world
// Ref: https://github.com/nodejs/node/issues/49344
const FinalizationRegistry = process.env.NODE_V8_COVERAGE ? FakeFinalizationRegistry : global.FinalizationRegistry || FakeFinalizationRegistry
const WeakRef = process.env.NODE_V8_COVERAGE ? FakeWeakRef : global.WeakRef || FakeWeakRef
const registry = new FinalizationRegistry((worker) => {
if (worker.exited) {
return
}
worker.terminate()
})
function createWorker (stream, opts) {
const { filename, workerData } = opts
const bundlerOverrides = '__bundlerPathsOverrides' in globalThis ? globalThis.__bundlerPathsOverrides : {}
const toExecute = bundlerOverrides['thread-stream-worker'] || join(__dirname, 'lib', 'worker.js')
const worker = new Worker(toExecute, {
...opts.workerOpts,
trackUnmanagedFds: false,
workerData: {
filename: filename.indexOf('file://') === 0
? filename
: pathToFileURL(filename).href,
dataBuf: stream[kImpl].dataBuf,
stateBuf: stream[kImpl].stateBuf,
workerData: {
$context: {
threadStreamVersion: version
},
...workerData
}
}
})
// We keep a strong reference for now,
// we need to start writing first
worker.stream = new FakeWeakRef(stream)
worker.on('message', onWorkerMessage)
worker.on('exit', onWorkerExit)
registry.register(stream, worker)
return worker
}
function drain (stream) {
assert(!stream[kImpl].sync)
if (stream[kImpl].needDrain) {
stream[kImpl].needDrain = false
stream.emit('drain')
}
}
function nextFlush (stream) {
const writeIndex = Atomics.load(stream[kImpl].state, WRITE_INDEX)
let leftover = stream[kImpl].data.length - writeIndex
if (leftover > 0) {
if (stream[kImpl].buf.length === 0) {
stream[kImpl].flushing = false
if (stream[kImpl].ending) {
end(stream)
} else if (stream[kImpl].needDrain) {
process.nextTick(drain, stream)
}
return
}
let toWrite = stream[kImpl].buf.slice(0, leftover)
let toWriteBytes = Buffer.byteLength(toWrite)
if (toWriteBytes <= leftover) {
stream[kImpl].buf = stream[kImpl].buf.slice(leftover)
// process._rawDebug('writing ' + toWrite.length)
write(stream, toWrite, nextFlush.bind(null, stream))
} else {
// multi-byte utf-8
stream.flush(() => {
// err is already handled in flush()
if (stream.destroyed) {
return
}
Atomics.store(stream[kImpl].state, READ_INDEX, 0)
Atomics.store(stream[kImpl].state, WRITE_INDEX, 0)
Atomics.notify(stream[kImpl].state, READ_INDEX)
// Find a toWrite length that fits the buffer
// it must exists as the buffer is at least 4 bytes length
// and the max utf-8 length for a char is 4 bytes.
while (toWriteBytes > stream[kImpl].data.length) {
leftover = leftover / 2
toWrite = stream[kImpl].buf.slice(0, leftover)
toWriteBytes = Buffer.byteLength(toWrite)
}
stream[kImpl].buf = stream[kImpl].buf.slice(leftover)
write(stream, toWrite, nextFlush.bind(null, stream))
})
}
} else if (leftover === 0) {
if (writeIndex === 0 && stream[kImpl].buf.length === 0) {
// we had a flushSync in the meanwhile
return
}
stream.flush(() => {
Atomics.store(stream[kImpl].state, READ_INDEX, 0)
Atomics.store(stream[kImpl].state, WRITE_INDEX, 0)
Atomics.notify(stream[kImpl].state, READ_INDEX)
nextFlush(stream)
})
} else {
// This should never happen
destroy(stream, new Error('overwritten'))
}
}
function onWorkerMessage (msg) {
const stream = this.stream.deref()
if (stream === undefined) {
this.exited = true
// Terminate the worker.
this.terminate()
return
}
switch (msg.code) {
case 'READY':
// Replace the FakeWeakRef with a
// proper one.
this.stream = new WeakRef(stream)
stream.flush(() => {
stream[kImpl].ready = true
stream.emit('ready')
})
break
case 'ERROR':
destroy(stream, msg.err)
break
case 'EVENT':
if (Array.isArray(msg.args)) {
stream.emit(msg.name, ...msg.args)
} else {
stream.emit(msg.name, msg.args)
}
break
case 'WARNING':
process.emitWarning(msg.err)
break
default:
destroy(stream, new Error('this should not happen: ' + msg.code))
}
}
function onWorkerExit (code) {
const stream = this.stream.deref()
if (stream === undefined) {
// Nothing to do, the worker already exit
return
}
registry.unregister(stream)
stream.worker.exited = true
stream.worker.off('exit', onWorkerExit)
destroy(stream, code !== 0 ? new Error('the worker thread exited') : null)
}
class ThreadStream extends EventEmitter {
constructor (opts = {}) {
super()
if (opts.bufferSize < 4) {
throw new Error('bufferSize must at least fit a 4-byte utf-8 char')
}
this[kImpl] = {}
this[kImpl].stateBuf = new SharedArrayBuffer(128)
this[kImpl].state = new Int32Array(this[kImpl].stateBuf)
this[kImpl].dataBuf = new SharedArrayBuffer(opts.bufferSize || 4 * 1024 * 1024)
this[kImpl].data = Buffer.from(this[kImpl].dataBuf)
this[kImpl].sync = opts.sync || false
this[kImpl].ending = false
this[kImpl].ended = false
this[kImpl].needDrain = false
this[kImpl].destroyed = false
this[kImpl].flushing = false
this[kImpl].ready = false
this[kImpl].finished = false
this[kImpl].errored = null
this[kImpl].closed = false
this[kImpl].buf = ''
// TODO (fix): Make private?
this.worker = createWorker(this, opts) // TODO (fix): make private
this.on('message', (message, transferList) => {
this.worker.postMessage(message, transferList)
})
}
write (data) {
if (this[kImpl].destroyed) {
error(this, new Error('the worker has exited'))
return false
}
if (this[kImpl].ending) {
error(this, new Error('the worker is ending'))
return false
}
if (this[kImpl].flushing && this[kImpl].buf.length + data.length >= MAX_STRING) {
try {
writeSync(this)
this[kImpl].flushing = true
} catch (err) {
destroy(this, err)
return false
}
}
this[kImpl].buf += data
if (this[kImpl].sync) {
try {
writeSync(this)
return true
} catch (err) {
destroy(this, err)
return false
}
}
if (!this[kImpl].flushing) {
this[kImpl].flushing = true
setImmediate(nextFlush, this)
}
this[kImpl].needDrain = this[kImpl].data.length - this[kImpl].buf.length - Atomics.load(this[kImpl].state, WRITE_INDEX) <= 0
return !this[kImpl].needDrain
}
end () {
if (this[kImpl].destroyed) {
return
}
this[kImpl].ending = true
end(this)
}
flush (cb) {
if (this[kImpl].destroyed) {
if (typeof cb === 'function') {
process.nextTick(cb, new Error('the worker has exited'))
}
return
}
// TODO write all .buf
const writeIndex = Atomics.load(this[kImpl].state, WRITE_INDEX)
// process._rawDebug(`(flush) readIndex (${Atomics.load(this.state, READ_INDEX)}) writeIndex (${Atomics.load(this.state, WRITE_INDEX)})`)
wait(this[kImpl].state, READ_INDEX, writeIndex, Infinity, (err, res) => {
if (err) {
destroy(this, err)
process.nextTick(cb, err)
return
}
if (res === 'not-equal') {
// TODO handle deadlock
this.flush(cb)
return
}
process.nextTick(cb)
})
}
flushSync () {
if (this[kImpl].destroyed) {
return
}
writeSync(this)
flushSync(this)
}
unref () {
this.worker.unref()
}
ref () {
this.worker.ref()
}
get ready () {
return this[kImpl].ready
}
get destroyed () {
return this[kImpl].destroyed
}
get closed () {
return this[kImpl].closed
}
get writable () {
return !this[kImpl].destroyed && !this[kImpl].ending
}
get writableEnded () {
return this[kImpl].ending
}
get writableFinished () {
return this[kImpl].finished
}
get writableNeedDrain () {
return this[kImpl].needDrain
}
get writableObjectMode () {
return false
}
get writableErrored () {
return this[kImpl].errored
}
}
function error (stream, err) {
setImmediate(() => {
stream.emit('error', err)
})
}
function destroy (stream, err) {
if (stream[kImpl].destroyed) {
return
}
stream[kImpl].destroyed = true
if (err) {
stream[kImpl].errored = err
error(stream, err)
}
if (!stream.worker.exited) {
stream.worker.terminate()
.catch(() => {})
.then(() => {
stream[kImpl].closed = true
stream.emit('close')
})
} else {
setImmediate(() => {
stream[kImpl].closed = true
stream.emit('close')
})
}
}
function write (stream, data, cb) {
// data is smaller than the shared buffer length
const current = Atomics.load(stream[kImpl].state, WRITE_INDEX)
const length = Buffer.byteLength(data)
stream[kImpl].data.write(data, current)
Atomics.store(stream[kImpl].state, WRITE_INDEX, current + length)
Atomics.notify(stream[kImpl].state, WRITE_INDEX)
cb()
return true
}
function end (stream) {
if (stream[kImpl].ended || !stream[kImpl].ending || stream[kImpl].flushing) {
return
}
stream[kImpl].ended = true
try {
stream.flushSync()
let readIndex = Atomics.load(stream[kImpl].state, READ_INDEX)
// process._rawDebug('writing index')
Atomics.store(stream[kImpl].state, WRITE_INDEX, -1)
// process._rawDebug(`(end) readIndex (${Atomics.load(stream.state, READ_INDEX)}) writeIndex (${Atomics.load(stream.state, WRITE_INDEX)})`)
Atomics.notify(stream[kImpl].state, WRITE_INDEX)
// Wait for the process to complete
let spins = 0
while (readIndex !== -1) {
// process._rawDebug(`read = ${read}`)
Atomics.wait(stream[kImpl].state, READ_INDEX, readIndex, 1000)
readIndex = Atomics.load(stream[kImpl].state, READ_INDEX)
if (readIndex === -2) {
destroy(stream, new Error('end() failed'))
return
}
if (++spins === 10) {
destroy(stream, new Error('end() took too long (10s)'))
return
}
}
process.nextTick(() => {
stream[kImpl].finished = true
stream.emit('finish')
})
} catch (err) {
destroy(stream, err)
}
// process._rawDebug('end finished...')
}
function writeSync (stream) {
const cb = () => {
if (stream[kImpl].ending) {
end(stream)
} else if (stream[kImpl].needDrain) {
process.nextTick(drain, stream)
}
}
stream[kImpl].flushing = false
while (stream[kImpl].buf.length !== 0) {
const writeIndex = Atomics.load(stream[kImpl].state, WRITE_INDEX)
let leftover = stream[kImpl].data.length - writeIndex
if (leftover === 0) {
flushSync(stream)
Atomics.store(stream[kImpl].state, READ_INDEX, 0)
Atomics.store(stream[kImpl].state, WRITE_INDEX, 0)
Atomics.notify(stream[kImpl].state, READ_INDEX)
continue
} else if (leftover < 0) {
// stream should never happen
throw new Error('overwritten')
}
let toWrite = stream[kImpl].buf.slice(0, leftover)
let toWriteBytes = Buffer.byteLength(toWrite)
if (toWriteBytes <= leftover) {
stream[kImpl].buf = stream[kImpl].buf.slice(leftover)
// process._rawDebug('writing ' + toWrite.length)
write(stream, toWrite, cb)
} else {
// multi-byte utf-8
flushSync(stream)
Atomics.store(stream[kImpl].state, READ_INDEX, 0)
Atomics.store(stream[kImpl].state, WRITE_INDEX, 0)
Atomics.notify(stream[kImpl].state, READ_INDEX)
// Find a toWrite length that fits the buffer
// it must exists as the buffer is at least 4 bytes length
// and the max utf-8 length for a char is 4 bytes.
while (toWriteBytes > stream[kImpl].buf.length) {
leftover = leftover / 2
toWrite = stream[kImpl].buf.slice(0, leftover)
toWriteBytes = Buffer.byteLength(toWrite)
}
stream[kImpl].buf = stream[kImpl].buf.slice(leftover)
write(stream, toWrite, cb)
}
}
}
function flushSync (stream) {
if (stream[kImpl].flushing) {
throw new Error('unable to flush while flushing')
}
// process._rawDebug('flushSync started')
const writeIndex = Atomics.load(stream[kImpl].state, WRITE_INDEX)
let spins = 0
// TODO handle deadlock
while (true) {
const readIndex = Atomics.load(stream[kImpl].state, READ_INDEX)
if (readIndex === -2) {
throw Error('_flushSync failed')
}
// process._rawDebug(`(flushSync) readIndex (${readIndex}) writeIndex (${writeIndex})`)
if (readIndex !== writeIndex) {
// TODO stream timeouts for some reason.
Atomics.wait(stream[kImpl].state, READ_INDEX, readIndex, 1000)
} else {
break
}
if (++spins === 10) {
throw new Error('_flushSync took too long (10s)')
}
}
// process._rawDebug('flushSync finished')
}
module.exports = ThreadStream
@@ -0,0 +1,9 @@
'use strict'
const WRITE_INDEX = 4
const READ_INDEX = 8
module.exports = {
WRITE_INDEX,
READ_INDEX
}
@@ -0,0 +1,68 @@
'use strict'
// Maximum wait time for a single waitAsync call
// Used as a fallback poll interval in case notifications are missed
// Keep this low enough for good throughput but high enough to not busy-loop
const WAIT_MS = 10000
function wait (state, index, expected, timeout, done) {
const max = timeout === Infinity ? Infinity : Date.now() + timeout
const check = () => {
const current = Atomics.load(state, index)
if (current === expected) {
done(null, 'ok')
return
}
if (max !== Infinity && Date.now() > max) {
done(null, 'timed-out')
return
}
// Wait for any change from current value
const remaining = max === Infinity ? WAIT_MS : Math.min(WAIT_MS, Math.max(1, max - Date.now()))
const result = Atomics.waitAsync(state, index, current, remaining)
if (result.async) {
result.value.then(check)
} else {
// Value already changed (not-equal) - recheck on next tick
setImmediate(check)
}
}
check()
}
function waitDiff (state, index, expected, timeout, done) {
const max = timeout === Infinity ? Infinity : Date.now() + timeout
const check = () => {
const current = Atomics.load(state, index)
if (current !== expected) {
done(null, 'ok')
return
}
if (max !== Infinity && Date.now() > max) {
done(null, 'timed-out')
return
}
// Wait for value to change from expected
const remaining = max === Infinity ? WAIT_MS : Math.min(WAIT_MS, Math.max(1, max - Date.now()))
const result = Atomics.waitAsync(state, index, expected, remaining)
if (result.async) {
result.value.then(check)
} else {
// Value already changed (not-equal) - recheck on next tick
setImmediate(check)
}
}
check()
}
module.exports = { wait, waitDiff }
@@ -0,0 +1,179 @@
'use strict'
const { realImport, realRequire } = require('real-require')
const { workerData, parentPort } = require('worker_threads')
const { WRITE_INDEX, READ_INDEX } = require('./indexes')
const { waitDiff } = require('./wait')
const {
dataBuf,
filename,
stateBuf
} = workerData
let destination
const state = new Int32Array(stateBuf)
const data = Buffer.from(dataBuf)
// Keep the event loop alive - Atomics.waitAsync promises don't prevent worker exit
const keepAlive = setInterval(() => {}, 60 * 60 * 1000)
async function start () {
let worker
try {
worker = (await realImport(filename))
} catch (error) {
// A yarn user that tries to start a ThreadStream for an external module
// provides a filename pointing to a zip file.
// eg. require.resolve('pino-elasticsearch') // returns /foo/pino-elasticsearch-npm-6.1.0-0c03079478-6915435172.zip/bar.js
// The `import` will fail to try to load it.
// This catch block executes the `require` fallback to load the module correctly.
// In fact, yarn modifies the `require` function to manage the zipped path.
// More details at https://github.com/pinojs/pino/pull/1113
// The error codes may change based on the node.js version (ENOTDIR > 12, ERR_MODULE_NOT_FOUND <= 12 )
if ((error.code === 'ENOTDIR' || error.code === 'ERR_MODULE_NOT_FOUND') &&
filename.startsWith('file://')) {
worker = realRequire(decodeURIComponent(filename.replace('file://', '')))
} else if (error.code === undefined || error.code === 'ERR_VM_DYNAMIC_IMPORT_CALLBACK_MISSING') {
// When bundled with pkg, an undefined error is thrown when called with realImport
// When bundled with pkg and using node v20, an ERR_VM_DYNAMIC_IMPORT_CALLBACK_MISSING error is thrown when called with realImport
// More info at: https://github.com/pinojs/thread-stream/issues/143
try {
worker = realRequire(decodeURIComponent(filename.replace(process.platform === 'win32' ? 'file:///' : 'file://', '')))
} catch {
throw error
}
} else if (filename.endsWith('.ts') || filename.endsWith('.cts')) {
// Native TypeScript import failed (type stripping not enabled).
// Fall back to ts-node for TypeScript files.
try {
if (!process[Symbol.for('ts-node.register.instance')]) {
realRequire('ts-node/register')
} else if (process.env.TS_NODE_DEV) {
realRequire('ts-node-dev')
}
worker = realRequire(decodeURIComponent(filename.replace(process.platform === 'win32' ? 'file:///' : 'file://', '')))
} catch {
throw error
}
} else {
throw error
}
}
// Depending on how the default export is performed, and on how the code is
// transpiled, we may find cases of two nested "default" objects.
// See https://github.com/pinojs/pino/issues/1243#issuecomment-982774762
if (typeof worker === 'object') worker = worker.default
if (typeof worker === 'object') worker = worker.default
destination = await worker(workerData.workerData)
destination.on('error', function (err) {
Atomics.store(state, WRITE_INDEX, -2)
Atomics.notify(state, WRITE_INDEX)
Atomics.store(state, READ_INDEX, -2)
Atomics.notify(state, READ_INDEX)
parentPort.postMessage({
code: 'ERROR',
err
})
})
destination.on('close', function () {
// process._rawDebug('worker close emitted')
const end = Atomics.load(state, WRITE_INDEX)
Atomics.store(state, READ_INDEX, end)
Atomics.notify(state, READ_INDEX)
clearInterval(keepAlive)
setImmediate(() => {
process.exit(0)
})
})
}
// No .catch() handler,
// in case there is an error it goes
// to unhandledRejection
start().then(function () {
parentPort.postMessage({
code: 'READY'
})
process.nextTick(run)
})
function run () {
const current = Atomics.load(state, READ_INDEX)
const end = Atomics.load(state, WRITE_INDEX)
// process._rawDebug(`pre state ${current} ${end}`)
if (end === current) {
if (end === data.length) {
waitDiff(state, READ_INDEX, end, Infinity, run)
} else {
waitDiff(state, WRITE_INDEX, end, Infinity, run)
}
return
}
// process._rawDebug(`post state ${current} ${end}`)
if (end === -1) {
// process._rawDebug('end')
destination.end()
return
}
const toWrite = data.toString('utf8', current, end)
// process._rawDebug('worker writing: ' + toWrite)
const res = destination.write(toWrite)
if (res) {
Atomics.store(state, READ_INDEX, end)
Atomics.notify(state, READ_INDEX)
setImmediate(run)
} else {
destination.once('drain', function () {
Atomics.store(state, READ_INDEX, end)
Atomics.notify(state, READ_INDEX)
run()
})
}
}
process.on('unhandledRejection', function (err) {
parentPort.postMessage({
code: 'ERROR',
err
})
process.exit(1)
})
process.on('uncaughtException', function (err) {
parentPort.postMessage({
code: 'ERROR',
err
})
process.exit(1)
})
process.once('exit', exitCode => {
if (exitCode !== 0) {
process.exit(exitCode)
return
}
if (destination?.writableNeedDrain && !destination?.writableEnded) {
parentPort.postMessage({
code: 'WARNING',
err: new Error('ThreadStream: process exited before destination stream was drained. this may indicate that the destination stream try to write to a another missing stream')
})
}
process.exit(0)
})
@@ -0,0 +1,52 @@
{
"name": "thread-stream",
"version": "4.0.0",
"description": "A streaming way to send data to a Node.js Worker Thread",
"main": "index.js",
"types": "index.d.ts",
"engines": {
"node": ">=20"
},
"dependencies": {
"real-require": "^0.2.0"
},
"devDependencies": {
"@types/node": "^22.0.0",
"@yao-pkg/pkg": "^6.0.0",
"borp": "^0.21.0",
"desm": "^1.3.0",
"eslint": "^9.39.1",
"fastbench": "^1.0.1",
"husky": "^9.0.6",
"neostandard": "^0.12.2",
"pino-elasticsearch": "^8.0.0",
"sonic-boom": "^4.0.1",
"ts-node": "^10.8.0",
"typescript": "~5.7.3"
},
"scripts": {
"build": "tsc --noEmit",
"lint": "eslint",
"test": "npm run lint && npm run build && npm run transpile && borp --pattern 'test/*.test.{js,mjs}'",
"test:ci": "npm run lint && npm run transpile && borp --pattern 'test/*.test.{js,mjs}'",
"test:yarn": "npm run transpile && borp --pattern 'test/*.test.js'",
"transpile": "sh ./test/ts/transpile.sh",
"prepare": "husky install"
},
"repository": {
"type": "git",
"url": "git+https://github.com/mcollina/thread-stream.git"
},
"keywords": [
"worker",
"thread",
"threads",
"stream"
],
"author": "Matteo Collina <hello@matteocollina.com>",
"license": "MIT",
"bugs": {
"url": "https://github.com/mcollina/thread-stream/issues"
},
"homepage": "https://github.com/mcollina/thread-stream#readme"
}
@@ -0,0 +1,259 @@
'use strict'
const { test } = require('node:test')
const assert = require('node:assert')
const { join } = require('path')
const { readFile } = require('fs')
const { file } = require('./helper')
const ThreadStream = require('..')
const { MessageChannel } = require('worker_threads')
const { once } = require('events')
test('base sync=true', function (t, done) {
const dest = file()
const stream = new ThreadStream({
filename: join(__dirname, 'to-file.js'),
workerData: { dest },
sync: true
})
assert.deepStrictEqual(stream.writableObjectMode, false)
assert.deepStrictEqual(stream.writableFinished, false)
stream.on('finish', () => {
assert.deepStrictEqual(stream.writableFinished, true)
readFile(dest, 'utf8', (err, data) => {
assert.ifError(err)
assert.strictEqual(data, 'hello world\nsomething else\n')
})
})
assert.deepStrictEqual(stream.closed, false)
stream.on('close', () => {
assert.deepStrictEqual(stream.closed, true)
assert.ok(!stream.writable)
done()
})
assert.deepStrictEqual(stream.writableNeedDrain, false)
assert.ok(stream.write('hello world\n'))
assert.ok(stream.write('something else\n'))
assert.ok(stream.writable)
assert.deepStrictEqual(stream.writableEnded, false)
stream.end()
assert.deepStrictEqual(stream.writableEnded, true)
})
test('overflow sync=true', function (t, done) {
const dest = file()
const stream = new ThreadStream({
bufferSize: 128,
filename: join(__dirname, 'to-file.js'),
workerData: { dest },
sync: true
})
let count = 0
// Write 10 chars, 20 times
function write () {
if (count++ === 20) {
stream.end()
return
}
stream.write('aaaaaaaaaa')
// do not wait for drain event
setImmediate(write)
}
write()
stream.on('close', () => {
readFile(dest, 'utf8', (err, data) => {
assert.ifError(err)
assert.strictEqual(data.length, 200)
done()
})
})
})
test('overflow sync=false', function (t, done) {
const dest = file()
const stream = new ThreadStream({
bufferSize: 128,
filename: join(__dirname, 'to-file.js'),
workerData: { dest },
sync: false
})
let count = 0
assert.deepStrictEqual(stream.writableNeedDrain, false)
// Write 10 chars, 20 times
function write () {
if (count++ === 20) {
stream.end()
return
}
if (!stream.write('aaaaaaaaaa')) {
assert.deepStrictEqual(stream.writableNeedDrain, true)
}
// do not wait for drain event
setImmediate(write)
}
write()
stream.on('drain', () => {
assert.deepStrictEqual(stream.writableNeedDrain, false)
})
stream.on('close', () => {
readFile(dest, 'utf8', (err, data) => {
assert.ifError(err)
assert.strictEqual(data.length, 200)
done()
})
})
})
test('over the bufferSize at startup', function (t, done) {
const dest = file()
const stream = new ThreadStream({
bufferSize: 10,
filename: join(__dirname, 'to-file.js'),
workerData: { dest },
sync: true
})
stream.on('finish', () => {
readFile(dest, 'utf8', (err, data) => {
assert.ifError(err)
assert.strictEqual(data, 'hello world\nsomething else\n')
})
})
stream.on('close', () => {
done()
})
assert.ok(stream.write('hello'))
assert.ok(stream.write(' world\n'))
assert.ok(stream.write('something else\n'))
stream.end()
})
test('over the bufferSize at startup (async)', function (t, done) {
const dest = file()
const stream = new ThreadStream({
bufferSize: 10,
filename: join(__dirname, 'to-file.js'),
workerData: { dest },
sync: false
})
assert.ok(stream.write('hello'))
assert.ok(!stream.write(' world\n'))
assert.ok(!stream.write('something else\n'))
stream.end()
stream.on('finish', () => {
readFile(dest, 'utf8', (err, data) => {
assert.ifError(err)
assert.strictEqual(data, 'hello world\nsomething else\n')
})
})
stream.on('close', () => {
done()
})
})
test('flushSync sync=false', function (t, done) {
const dest = file()
const stream = new ThreadStream({
bufferSize: 128,
filename: join(__dirname, 'to-file.js'),
workerData: { dest },
sync: false
})
stream.on('drain', () => {
stream.end()
})
stream.on('close', () => {
readFile(dest, 'utf8', (err, data) => {
assert.ifError(err)
assert.strictEqual(data.length, 200)
done()
})
})
for (let count = 0; count < 20; count++) {
stream.write('aaaaaaaaaa')
}
stream.flushSync()
})
test('pass down MessagePorts', async function (t) {
const { port1, port2 } = new MessageChannel()
const stream = new ThreadStream({
filename: join(__dirname, 'port.js'),
workerData: { port: port1 },
workerOpts: {
transferList: [port1]
},
sync: false
})
t.after(() => {
stream.end()
})
assert.ok(stream.write('hello world\n'))
assert.ok(stream.write('something else\n'))
const [strings] = await once(port2, 'message')
assert.strictEqual(strings, 'hello world\nsomething else\n')
})
test('destroy does not error', function (t, done) {
const dest = file()
const stream = new ThreadStream({
filename: join(__dirname, 'to-file.js'),
workerData: { dest },
sync: false
})
stream.on('ready', () => {
stream.worker.terminate()
})
stream.on('error', (err) => {
assert.strictEqual(err.message, 'the worker thread exited')
stream.flush((err) => {
assert.strictEqual(err.message, 'the worker has exited')
})
assert.doesNotThrow(() => stream.flushSync())
assert.doesNotThrow(() => stream.end())
done()
})
})
test('syntax error', function (t, done) {
const stream = new ThreadStream({
filename: join(__dirname, 'syntax-error.mjs')
})
stream.on('error', (err) => {
assert.strictEqual(err.message, 'Unexpected end of input')
done()
})
})
@@ -0,0 +1,38 @@
'use strict'
const { test } = require('node:test')
const { join } = require('path')
const ThreadStream = require('..')
const { file } = require('./helper')
const MAX = 1000
let str = ''
for (let i = 0; i < 10; i++) {
str += 'hello'
}
test('base', function (t, done) {
const dest = file()
const stream = new ThreadStream({
filename: join(__dirname, 'to-file.js'),
workerData: { dest }
})
let runs = 0
function benchThreadStream () {
if (++runs === 1000) {
stream.end()
return
}
for (let i = 0; i < MAX; i++) {
stream.write(str)
}
setImmediate(benchThreadStream)
}
benchThreadStream()
stream.on('finish', function () {
done()
})
})
@@ -0,0 +1,59 @@
'use strict'
const { test } = require('node:test')
const assert = require('node:assert')
const { join } = require('path')
const { file } = require('./helper')
const ThreadStream = require('..')
test('bundlers support with .js file', function (t, done) {
globalThis.__bundlerPathsOverrides = {
'thread-stream-worker': join(__dirname, 'custom-worker.js')
}
const dest = file()
process.on('uncaughtException', error => {
console.log(error)
})
const stream = new ThreadStream({
filename: join(__dirname, 'to-file.js'),
workerData: { dest },
sync: true
})
stream.worker.removeAllListeners('message')
stream.worker.once('message', message => {
assert.strictEqual(message.code, 'CUSTOM-WORKER-CALLED')
done()
})
stream.end()
})
test('bundlers support with .mjs file', function (t, done) {
globalThis.__bundlerPathsOverrides = {
'thread-stream-worker': join(__dirname, 'custom-worker.js')
}
const dest = file()
process.on('uncaughtException', error => {
console.log(error)
})
const stream = new ThreadStream({
filename: join(__dirname, 'to-file.mjs'),
workerData: { dest },
sync: true
})
stream.worker.removeAllListeners('message')
stream.worker.once('message', message => {
assert.strictEqual(message.code, 'CUSTOM-WORKER-CALLED')
done()
})
stream.end()
})
@@ -0,0 +1,37 @@
'use strict'
const { join } = require('path')
const ThreadStream = require('..')
const assert = require('assert')
let worker = null
function setup () {
const stream = new ThreadStream({
filename: join(__dirname, 'to-file.js'),
workerData: { dest: process.argv[2] },
sync: true
})
worker = stream.worker
stream.write('hello')
stream.write(' ')
stream.write('world\n')
stream.flushSync()
stream.unref()
// the stream object goes out of scope here
setImmediate(gc) // eslint-disable-line
}
setup()
let exitEmitted = false
worker.on('exit', function () {
exitEmitted = true
})
process.on('exit', function () {
assert.strictEqual(exitEmitted, true)
})
@@ -0,0 +1,75 @@
'use strict'
const { test } = require('node:test')
const assert = require('node:assert')
const { join } = require('path')
const { MessageChannel } = require('worker_threads')
const { once } = require('events')
const ThreadStream = require('..')
const isYarnPnp = process.versions.pnp !== undefined
test('yarn module resolution', { skip: !isYarnPnp }, (t, done) => {
const modulePath = require.resolve('pino-elasticsearch')
assert.match(modulePath, /.*\.zip.*/)
const stream = new ThreadStream({
filename: modulePath,
workerData: { node: null },
sync: true
})
assert.deepStrictEqual(stream.writableErrored, null)
stream.on('error', (err) => {
assert.deepStrictEqual(stream.writableErrored, err)
})
assert.ok(stream.write('hello world\n'))
assert.ok(stream.writable)
stream.end()
done()
})
test('yarn module resolution for directories with special characters', { skip: !isYarnPnp }, async t => {
const { port1, port2 } = new MessageChannel()
const stream = new ThreadStream({
filename: join(__dirname, 'dir with spaces', 'test-package.zip', 'worker.js'),
workerData: { port: port1 },
workerOpts: {
transferList: [port1]
},
sync: false
})
t.after(() => {
stream.end()
})
assert.ok(stream.write('hello world\n'))
assert.ok(stream.write('something else\n'))
const [strings] = await once(port2, 'message')
assert.strictEqual(strings, 'hello world\nsomething else\n')
})
test('yarn module resolution for typescript commonjs modules', { skip: !isYarnPnp }, async t => {
const { port1, port2 } = new MessageChannel()
const stream = new ThreadStream({
filename: join(__dirname, 'ts-commonjs-default-export.zip', 'worker.js'),
workerData: { port: port1 },
workerOpts: {
transferList: [port1]
},
sync: false
})
t.after(() => {
stream.end()
})
assert.ok(stream.write('hello world\n'))
assert.ok(stream.write('something else\n'))
const [strings] = await once(port2, 'message')
assert.strictEqual(strings, 'hello world\nsomething else\n')
})
@@ -0,0 +1,21 @@
'use strict'
const { test } = require('node:test')
const assert = require('node:assert')
const { join } = require('path')
const ThreadStream = require('..')
const { version } = require('../package.json')
test('get context', (t, done) => {
const stream = new ThreadStream({
filename: join(__dirname, 'get-context.js'),
workerData: {},
sync: true
})
t.after(() => stream.end())
stream.on('context', (ctx) => {
assert.deepStrictEqual(ctx.threadStreamVersion, version)
done()
})
stream.write('hello')
})
@@ -0,0 +1,16 @@
'use strict'
const { join } = require('path')
const ThreadStream = require('..')
const stream = new ThreadStream({
filename: join(__dirname, 'to-file.js'),
workerData: { dest: process.argv[2] },
sync: true
})
stream.write('hello')
stream.write(' ')
stream.write('world\n')
stream.flushSync()
stream.unref()
@@ -0,0 +1,9 @@
'use strict'
const { parentPort } = require('worker_threads')
parentPort.postMessage({
code: 'CUSTOM-WORKER-CALLED'
})
require('../lib/worker')
@@ -0,0 +1,22 @@
'use strict'
const { Writable } = require('stream')
const parentPort = require('worker_threads').parentPort
async function run () {
return new Writable({
autoDestroy: true,
write (chunk, enc, cb) {
if (parentPort) {
parentPort.postMessage({
code: 'EVENT',
name: 'socketError',
args: ['list', 'of', 'args', 123, new Error('unable to write data to the TCP socket')]
})
}
cb()
}
})
}
module.exports = run
@@ -0,0 +1,56 @@
'use strict'
const { test } = require('node:test')
const assert = require('node:assert')
const { join } = require('path')
const { readFile } = require('fs')
const { file } = require('./helper')
const ThreadStream = require('..')
test('destroy support', function (t, done) {
const dest = file()
const stream = new ThreadStream({
filename: join(__dirname, 'to-file-on-destroy.js'),
workerData: { dest },
sync: true
})
stream.on('close', () => {
assert.ok(!stream.writable)
readFile(dest, 'utf8', (err, data) => {
assert.ifError(err)
assert.strictEqual(data, 'hello world\nsomething else\n')
done()
})
})
assert.ok(stream.write('hello world\n'))
assert.ok(stream.write('something else\n'))
assert.ok(stream.writable)
stream.end()
})
test('synchronous _final support', function (t, done) {
const dest = file()
const stream = new ThreadStream({
filename: join(__dirname, 'to-file-on-final.js'),
workerData: { dest },
sync: true
})
stream.on('close', () => {
assert.ok(!stream.writable)
readFile(dest, 'utf8', (err, data) => {
assert.ifError(err)
assert.strictEqual(data, 'hello world\nsomething else\n')
done()
})
})
assert.ok(stream.write('hello world\n'))
assert.ok(stream.write('something else\n'))
assert.ok(stream.writable)
stream.end()
})
@@ -0,0 +1,14 @@
'use strict'
const { Writable } = require('stream')
async function run (opts) {
const stream = new Writable({
write (chunk, enc, cb) {
cb(new Error('kaboom'))
}
})
return stream
}
module.exports = run
@@ -0,0 +1,46 @@
import { test } from 'node:test'
import assert from 'node:assert'
import { readFile } from 'fs'
import ThreadStream from '../index.js'
import { join } from 'desm'
import { pathToFileURL } from 'url'
import { file } from './helper.js'
function basic (text, filename) {
test(text, function (t, done) {
const dest = file()
const stream = new ThreadStream({
filename,
workerData: { dest },
sync: true
})
stream.on('finish', () => {
readFile(dest, 'utf8', (err, data) => {
assert.ifError(err)
assert.strictEqual(data, 'hello world\nsomething else\n')
})
})
stream.on('close', () => {
done()
})
assert.ok(stream.write('hello world\n'))
assert.ok(stream.write('something else\n'))
stream.end()
})
}
basic('esm with path', join(import.meta.url, 'to-file.mjs'))
basic('esm with file URL', pathToFileURL(join(import.meta.url, 'to-file.mjs')).href)
basic('(ts -> es6) esm with path', join(import.meta.url, 'ts', 'to-file.es6.mjs'))
basic('(ts -> es6) esm with file URL', pathToFileURL(join(import.meta.url, 'ts', 'to-file.es6.mjs')).href)
basic('(ts -> es2017) esm with path', join(import.meta.url, 'ts', 'to-file.es2017.mjs'))
basic('(ts -> es2017) esm with file URL', pathToFileURL(join(import.meta.url, 'ts', 'to-file.es2017.mjs')).href)
basic('(ts -> esnext) esm with path', join(import.meta.url, 'ts', 'to-file.esnext.mjs'))
basic('(ts -> esnext) esm with file URL', pathToFileURL(join(import.meta.url, 'ts', 'to-file.esnext.mjs')).href)
@@ -0,0 +1,24 @@
'use strict'
const { test } = require('node:test')
const assert = require('node:assert')
const { join } = require('path')
const ThreadStream = require('..')
test('event propagate', (t, done) => {
const stream = new ThreadStream({
filename: join(__dirname, 'emit-event.js'),
workerData: {},
sync: true
})
t.after(() => stream.end())
stream.on('socketError', function (a, b, c, n, error) {
assert.deepStrictEqual(a, 'list')
assert.deepStrictEqual(b, 'of')
assert.deepStrictEqual(c, 'args')
assert.deepStrictEqual(n, 123)
assert.deepStrictEqual(error, new Error('unable to write data to the TCP socket'))
done()
})
stream.write('hello')
})
@@ -0,0 +1,14 @@
'use strict'
const { Writable } = require('stream')
async function run (opts) {
const stream = new Writable({
write (chunk, enc, cb) {
process.exit(1)
}
})
return stream
}
module.exports = run
@@ -0,0 +1,22 @@
'use strict'
const { Writable } = require('stream')
const parentPort = require('worker_threads').parentPort
async function run (opts) {
return new Writable({
autoDestroy: true,
write (chunk, enc, cb) {
if (parentPort) {
parentPort.postMessage({
code: 'EVENT',
name: 'context',
args: opts.$context
})
}
cb()
}
})
}
module.exports = run
@@ -0,0 +1,26 @@
'use strict'
const { join } = require('path')
const { tmpdir } = require('os')
const { unlinkSync } = require('fs')
const files = []
let count = 0
function file () {
const file = join(tmpdir(), `thread-stream-${process.pid}-${count++}`)
files.push(file)
return file
}
process.on('beforeExit', () => {
for (const file of files) {
try {
unlinkSync(file)
} catch (e) {
// ignore cleanup errors
}
}
})
module.exports.file = file
@@ -0,0 +1,11 @@
'use strict'
const { test } = require('node:test')
const assert = require('node:assert')
const indexes = require('../lib/indexes')
for (const index of Object.keys(indexes)) {
test(`${index} is lock free`, function () {
assert.strictEqual(Atomics.isLockFree(indexes[index]), true)
})
}
@@ -0,0 +1,75 @@
import { test } from 'node:test'
import assert from 'node:assert'
import { readFile } from 'fs'
import ThreadStream from '../index.js'
import { join } from 'desm'
import { file } from './helper.js'
test('break up utf8 multibyte (sync)', (t, done) => {
const longString = '\u03A3'.repeat(16)
const dest = file()
const stream = new ThreadStream({
bufferSize: 15, // this must be odd
filename: join(import.meta.url, 'to-file.js'),
workerData: { dest },
sync: true
})
stream.on('finish', () => {
readFile(dest, 'utf8', (err, data) => {
assert.ifError(err)
assert.strictEqual(data, longString)
done()
})
})
stream.write(longString)
stream.end()
})
test('break up utf8 multibyte (async)', (t, done) => {
const longString = '\u03A3'.repeat(16)
const dest = file()
const stream = new ThreadStream({
bufferSize: 15, // this must be odd
filename: join(import.meta.url, 'to-file.js'),
workerData: { dest },
sync: false
})
stream.on('finish', () => {
readFile(dest, 'utf8', (err, data) => {
assert.ifError(err)
assert.strictEqual(data, longString)
done()
})
})
stream.write(longString)
stream.end()
})
test('break up utf8 multibyte several times bigger than write buffer', (t, done) => {
const longString = '\u03A3'.repeat(32)
const dest = file()
const stream = new ThreadStream({
bufferSize: 15, // this must be odd
filename: join(import.meta.url, 'to-file.js'),
workerData: { dest },
sync: false
})
stream.on('finish', () => {
readFile(dest, 'utf8', (err, data) => {
assert.ifError(err)
assert.strictEqual(data, longString)
done()
})
})
stream.write(longString)
stream.end()
})
@@ -0,0 +1,18 @@
'use strict'
const { parentPort } = require('worker_threads')
const { Writable } = require('stream')
function run () {
parentPort.once('message', function ({ text, takeThisPortPlease }) {
takeThisPortPlease.postMessage(`received: ${text}`)
})
return new Writable({
autoDestroy: true,
write (chunk, enc, cb) {
cb()
}
})
}
module.exports = run
@@ -0,0 +1,37 @@
'use strict'
/**
* This file is packaged using pkg in order to test if worker.js works in that context.
* Note: We can't use node:test here because it crashes inside pkg bundles due to V8 internals.
*/
const assert = require('node:assert')
const { join } = require('path')
const { file } = require('../helper')
const ThreadStream = require('../..')
globalThis.__bundlerPathsOverrides = {
'thread-stream-worker': join(__dirname, '..', 'custom-worker.js')
}
const dest = file()
process.on('uncaughtException', (error) => {
console.error(error)
process.exit(1)
})
const stream = new ThreadStream({
filename: join(__dirname, '..', 'to-file.js'),
workerData: { dest },
sync: true
})
stream.worker.removeAllListeners('message')
stream.worker.once('message', (message) => {
assert.strictEqual(message.code, 'CUSTOM-WORKER-CALLED')
console.log('pkg test passed')
process.exit(0)
})
stream.end()
@@ -0,0 +1,14 @@
{
"pkg": {
"assets": [
"../custom-worker.js",
"../to-file.js"
],
"targets": [
"node20",
"node22",
"node24"
],
"outputPath": "test/pkg"
}
}
@@ -0,0 +1,45 @@
'use strict'
const { test } = require('node:test')
const assert = require('node:assert')
const config = require('./pkg.config.json')
const { promisify } = require('util')
const { unlink } = require('fs/promises')
const { join } = require('path')
const { platform } = require('process')
const exec = promisify(require('child_process').exec)
test('worker test when packaged into executable using pkg', async () => {
const packageName = 'index'
// package the app into several node versions, check config for more info
const filePath = `${join(__dirname, packageName)}.js`
const configPath = join(__dirname, 'pkg.config.json')
process.env.NODE_OPTIONS ||= ''
process.env.NODE_OPTIONS = '--no-warnings'
const { stderr } = await exec(`npx pkg ${filePath} --config ${configPath}`)
// there should be no error when packaging
assert.strictEqual(stderr, '')
// pkg outputs files in the following format by default: {filename}-{node version}
for (const target of config.pkg.targets) {
// execute the packaged test
let executablePath = `${join(config.pkg.outputPath, packageName)}-${target}`
// when on windows, we need the .exe extension
if (platform === 'win32') {
executablePath = `${executablePath}.exe`
} else {
executablePath = `./${executablePath}`
}
const { stderr } = await exec(executablePath)
// check if there were no errors
assert.strictEqual(stderr, '')
// clean up afterwards
await unlink(executablePath)
}
})
@@ -0,0 +1,16 @@
'use strict'
const { Writable } = require('stream')
function run (opts) {
const { port } = opts
return new Writable({
autoDestroy: true,
write (chunk, enc, cb) {
port.postMessage(chunk.toString())
cb()
}
})
}
module.exports = run
@@ -0,0 +1,23 @@
'use strict'
const { test } = require('node:test')
const assert = require('node:assert')
const { join } = require('path')
const { once } = require('events')
const { MessageChannel } = require('worker_threads')
const ThreadStream = require('..')
test('message events emitted on the stream are posted to the worker', async function (t) {
const { port1, port2 } = new MessageChannel()
const stream = new ThreadStream({
filename: join(__dirname, 'on-message.js'),
sync: false
})
t.after(() => {
stream.end()
})
stream.emit('message', { text: 'hello', takeThisPortPlease: port1 }, [port1])
const [confirmation] = await once(port2, 'message')
assert.strictEqual(confirmation, 'received: hello')
})
@@ -0,0 +1,35 @@
'use strict'
const { test } = require('node:test')
const assert = require('node:assert')
const { join } = require('path')
const { file } = require('./helper')
const { createReadStream } = require('fs')
const ThreadStream = require('..')
const buffer = require('buffer')
const MAX_STRING = buffer.constants.MAX_STRING_LENGTH
test('string limit 2', { skip: process.env.CI }, (t, done) => {
const dest = file()
const stream = new ThreadStream({
filename: join(__dirname, 'to-file.js'),
workerData: { dest },
sync: false
})
stream.on('close', async () => {
let buf
for await (const chunk of createReadStream(dest)) {
buf = chunk
}
assert.strictEqual('asd', buf.toString().slice(-3))
done()
})
stream.on('ready', () => {
stream.write('a'.repeat(MAX_STRING - 2))
stream.write('asd')
stream.end()
})
})
@@ -0,0 +1,37 @@
'use strict'
const { test } = require('node:test')
const assert = require('node:assert')
const { join } = require('path')
const { file } = require('./helper')
const { stat } = require('fs')
const ThreadStream = require('..')
test('string limit', { skip: process.env.CI, timeout: 30000 }, (t, done) => {
const dest = file()
const stream = new ThreadStream({
filename: join(__dirname, 'to-file.js'),
workerData: { dest },
sync: false
})
let length = 0
stream.on('close', () => {
stat(dest, (err, f) => {
assert.ifError(err)
assert.strictEqual(f.size, length)
done()
})
})
const buf = Buffer.alloc(1024).fill('x').toString() // 1 KB
// This writes 1 GB of data
for (let i = 0; i < 1024 * 1024; i++) {
length += buf.length
stream.write(buf)
}
stream.end()
})
@@ -0,0 +1,2 @@
// this is a syntax error
import
@@ -0,0 +1,138 @@
'use strict'
const { test } = require('node:test')
const assert = require('node:assert')
const { fork } = require('child_process')
const { join } = require('path')
const { readFile } = require('fs').promises
const { file } = require('./helper')
const { once } = require('events')
const ThreadStream = require('..')
test('exits with 0', async function (t) {
const dest = file()
const child = fork(join(__dirname, 'create-and-exit.js'), [dest])
const [code] = await once(child, 'exit')
assert.strictEqual(code, 0)
const data = await readFile(dest, 'utf8')
assert.strictEqual(data, 'hello world\n')
})
test('emit error if thread exits', async function (t) {
const stream = new ThreadStream({
filename: join(__dirname, 'exit.js'),
sync: true
})
const closed = once(stream, 'close').catch(() => {})
stream.on('ready', () => {
stream.write('hello world\n')
})
let [err] = await once(stream, 'error')
assert.strictEqual(err.message, 'the worker thread exited')
stream.write('noop');
[err] = await once(stream, 'error')
assert.strictEqual(err.message, 'the worker has exited')
stream.write('noop');
[err] = await once(stream, 'error')
assert.strictEqual(err.message, 'the worker has exited')
await closed
})
test('emit error if thread have unhandledRejection', async function (t) {
const stream = new ThreadStream({
filename: join(__dirname, 'unhandledRejection.js'),
sync: true
})
const closed = once(stream, 'close').catch(() => {})
stream.on('ready', () => {
stream.write('hello world\n')
})
let [err] = await once(stream, 'error')
assert.strictEqual(err.message, 'kaboom')
stream.write('noop');
[err] = await once(stream, 'error')
assert.strictEqual(err.message, 'the worker has exited')
stream.write('noop');
[err] = await once(stream, 'error')
assert.strictEqual(err.message, 'the worker has exited')
await closed
})
test('emit error if worker stream emit error', async function (t) {
const stream = new ThreadStream({
filename: join(__dirname, 'error.js'),
sync: true
})
const closed = once(stream, 'close').catch(() => {})
stream.on('ready', () => {
stream.write('hello world\n')
})
let [err] = await once(stream, 'error')
assert.strictEqual(err.message, 'kaboom')
stream.write('noop');
[err] = await once(stream, 'error')
assert.strictEqual(err.message, 'the worker has exited')
stream.write('noop');
[err] = await once(stream, 'error')
assert.strictEqual(err.message, 'the worker has exited')
await closed
})
test('emit error if thread have uncaughtException', async function (t) {
const stream = new ThreadStream({
filename: join(__dirname, 'uncaughtException.js'),
sync: true
})
const closed = once(stream, 'close').catch(() => {})
stream.on('ready', () => {
stream.write('hello world\n')
})
let [err] = await once(stream, 'error')
assert.strictEqual(err.message, 'kaboom')
stream.write('noop');
[err] = await once(stream, 'error')
assert.strictEqual(err.message, 'the worker has exited')
stream.write('noop');
[err] = await once(stream, 'error')
assert.strictEqual(err.message, 'the worker has exited')
await closed
})
test('close the work if out of scope on gc', { skip: !global.WeakRef }, async function (t) {
const dest = file()
const child = fork(join(__dirname, 'close-on-gc.js'), [dest], {
execArgv: ['--expose-gc']
})
const [code] = await once(child, 'exit')
assert.strictEqual(code, 0)
const data = await readFile(dest, 'utf8')
assert.strictEqual(data, 'hello world\n')
})
@@ -0,0 +1,23 @@
'use strict'
const fs = require('fs')
const { Writable } = require('stream')
function run (opts) {
let data = ''
return new Writable({
autoDestroy: true,
write (chunk, enc, cb) {
data += chunk.toString()
cb()
},
destroy (err, cb) {
// process._rawDebug('destroy called')
fs.writeFile(opts.dest, data, function (err2) {
cb(err2 || err)
})
}
})
}
module.exports = run
@@ -0,0 +1,24 @@
'use strict'
const fs = require('fs')
const { Writable } = require('stream')
function run (opts) {
let data = ''
return new Writable({
autoDestroy: true,
write (chunk, enc, cb) {
data += chunk.toString()
cb()
},
final (cb) {
setTimeout(function () {
fs.writeFile(opts.dest, data, function (err) {
cb(err)
})
}, 100)
}
})
}
module.exports = run
@@ -0,0 +1,12 @@
'use strict'
const fs = require('fs')
const { once } = require('events')
async function run (opts) {
const stream = fs.createWriteStream(opts.dest)
await once(stream, 'open')
return stream
}
module.exports = run
@@ -0,0 +1,8 @@
import { createWriteStream } from 'fs'
import { once } from 'events'
export default async function run (opts) {
const stream = createWriteStream(opts.dest)
await once(stream, 'open')
return stream
}
@@ -0,0 +1,9 @@
'use strict'
const { PassThrough } = require('stream')
async function run (opts) {
return new PassThrough({})
}
module.exports = run
@@ -0,0 +1,29 @@
'use strict'
const { test } = require('node:test')
const assert = require('node:assert')
const { join } = require('path')
const { file } = require('./helper')
const ThreadStream = require('..')
function basic (esVersion) {
test(`transpiled-ts-to-${esVersion}`, function () {
const dest = file()
const stream = new ThreadStream({
filename: join(__dirname, 'ts', `to-file.${esVersion}.cjs`),
workerData: { dest },
sync: true
})
// There are arbitrary checks, the important aspect of this test is to ensure
// that we can properly load the transpiled file into our worker thread.
assert.deepStrictEqual(stream.writableEnded, false)
stream.end()
assert.deepStrictEqual(stream.writableEnded, true)
})
}
basic('es5')
basic('es6')
basic('es2017')
basic('esnext')
@@ -0,0 +1,35 @@
import { test } from 'node:test'
import assert from 'node:assert'
import { readFile } from 'fs/promises'
import ThreadStream from '../index.js'
import { join } from 'desm'
import { file } from './helper.js'
const nodeVersion = parseInt(process.versions.node.split('.')[0], 10)
// Native TypeScript stripping (--experimental-strip-types) is only available in Node 22.6+
test('typescript module with native type stripping', { skip: nodeVersion < 22 }, async function (t) {
const dest = file()
const stream = new ThreadStream({
filename: join(import.meta.url, 'ts', 'to-file.ts'),
workerData: { dest },
workerOpts: {
execArgv: ['--experimental-strip-types', '--disable-warning=ExperimentalWarning']
},
sync: false
})
t.after(() => stream.end())
assert.ok(stream.write('hello world\n'))
assert.ok(stream.write('something else\n'))
stream.end()
await new Promise((resolve) => {
stream.on('close', resolve)
})
const data = await readFile(dest, 'utf8')
assert.strictEqual(data, 'hello world\nsomething else\n')
})
@@ -0,0 +1,35 @@
'use strict'
const { test } = require('node:test')
const assert = require('node:assert')
const { readFile } = require('fs/promises')
const { join } = require('path')
const { file } = require('./helper')
const ThreadStream = require('..')
// This test verifies that TypeScript files can be loaded via ts-node
// when native type stripping is not enabled in the worker thread.
// Unlike ts.test.ts which passes --experimental-strip-types via execArgv,
// this test does NOT pass that flag, so the worker will fall back to ts-node.
test('typescript module with ts-node fallback', async function (t) {
const dest = file()
const stream = new ThreadStream({
filename: join(__dirname, 'ts', 'to-file.ts'),
workerData: { dest },
sync: false
})
t.after(() => stream.end())
assert.ok(stream.write('hello world\n'))
assert.ok(stream.write('something else\n'))
stream.end()
await new Promise((resolve) => {
stream.on('close', resolve)
})
const data = await readFile(dest, 'utf8')
assert.strictEqual(data, 'hello world\nsomething else\n')
})
@@ -0,0 +1,10 @@
import { type PathLike, type WriteStream, createWriteStream } from 'fs'
import { once } from 'events'
export default async function run (
opts: { dest: PathLike },
): Promise<WriteStream> {
const stream = createWriteStream(opts.dest)
await once(stream, 'open')
return stream
}
@@ -0,0 +1,19 @@
#!/bin/sh
set -e
cd ./test/ts;
if (echo "${npm_config_user_agent}" | grep "yarn"); then
export RUNNER="yarn";
else
export RUNNER="npx";
fi
test ./to-file.ts -ot ./to-file.es5.cjs || ("${RUNNER}" tsc --skipLibCheck --target es5 ./to-file.ts && mv ./to-file.js ./to-file.es5.cjs);
test ./to-file.ts -ot ./to-file.es6.mjs || ("${RUNNER}" tsc --skipLibCheck --target es6 ./to-file.ts && mv ./to-file.js ./to-file.es6.mjs);
test ./to-file.ts -ot ./to-file.es6.cjs || ("${RUNNER}" tsc --skipLibCheck --target es6 --module commonjs ./to-file.ts && mv ./to-file.js ./to-file.es6.cjs);
test ./to-file.ts -ot ./to-file.es2017.mjs || ("${RUNNER}" tsc --skipLibCheck --target es2017 ./to-file.ts && mv ./to-file.js ./to-file.es2017.mjs);
test ./to-file.ts -ot ./to-file.es2017.cjs || ("${RUNNER}" tsc --skipLibCheck --target es2017 --module commonjs ./to-file.ts && mv ./to-file.js ./to-file.es2017.cjs);
test ./to-file.ts -ot ./to-file.esnext.mjs || ("${RUNNER}" tsc --skipLibCheck --target esnext --module esnext ./to-file.ts && mv ./to-file.js ./to-file.esnext.mjs);
test ./to-file.ts -ot ./to-file.esnext.cjs || ("${RUNNER}" tsc --skipLibCheck --target esnext --module commonjs ./to-file.ts && mv ./to-file.js ./to-file.esnext.cjs);
@@ -0,0 +1,21 @@
'use strict'
const { Writable } = require('stream')
// Nop console.error to avoid printing things out
console.error = () => {}
setImmediate(function () {
throw new Error('kaboom')
})
async function run (opts) {
const stream = new Writable({
write (chunk, enc, cb) {
cb()
}
})
return stream
}
module.exports = run
@@ -0,0 +1,21 @@
'use strict'
const { Writable } = require('stream')
// Nop console.error to avoid printing things out
console.error = () => {}
setImmediate(function () {
Promise.reject(new Error('kaboom'))
})
async function run (opts) {
const stream = new Writable({
write (chunk, enc, cb) {
cb()
}
})
return stream
}
module.exports = run
@@ -0,0 +1,7 @@
nodeLinker: pnp
pnpMode: loose
pnpEnableEsmLoader: false
packageExtensions:
debug@*:
dependencies:
supports-color: '*'
@@ -0,0 +1,8 @@
{
"compilerOptions": {
"esModuleInterop": true
},
"files": [
"index.d.ts"
],
}