diff --git a/doc/api/test.md b/doc/api/test.md index e5ce08a76303..8cdfa8845892 100644 --- a/doc/api/test.md +++ b/doc/api/test.md @@ -4575,7 +4575,9 @@ The unique identifier of the worker running the current test file. This value is derived from the `NODE_TEST_WORKER_ID` environment variable. When running tests with `--test-isolation=process` (the default), each test file runs in a separate child process and is assigned a worker ID from 1 to N, where N is the number of -concurrent workers. When running with `--test-isolation=none`, all tests run in +concurrent workers. A worker ID is never shared by two test files running at the +same time. Once a test file finishes, its worker ID is reused by the next test +file that starts. When running with `--test-isolation=none`, all tests run in the same process and the worker ID is always 1. This value is `undefined` when not running in a test context. diff --git a/lib/internal/test_runner/runner.js b/lib/internal/test_runner/runner.js index da3bd744ffc1..5f73e04a6366 100644 --- a/lib/internal/test_runner/runner.js +++ b/lib/internal/test_runner/runner.js @@ -15,7 +15,6 @@ const { ArrayPrototypeSlice, ArrayPrototypeSome, ArrayPrototypeSort, - MathMax, ObjectAssign, PromisePrototypeThen, PromiseWithResolvers, @@ -38,7 +37,6 @@ const { const { spawn } = require('child_process'); const { statSync } = require('fs'); const { finished } = require('internal/streams/end-of-stream'); -const { availableParallelism } = require('os'); const { resolve, sep, isAbsolute } = require('path'); const { DefaultDeserializer, DefaultSerializer } = require('v8'); const { getOptionValue, getOptionsAsFlagsFromBinding } = require('internal/options'); @@ -139,17 +137,22 @@ let kResistStopPropagation; // Worker ID pool management for concurrent test execution class WorkerIdPool { - #nextId = 0; - #maxConcurrency; - - constructor(maxConcurrency) { - this.#maxConcurrency = maxConcurrency; - } + #acquiredIds = new SafeSet(); acquire() { - const id = (this.#nextId++ % this.#maxConcurrency) + 1; + let id = 1; + + while (this.#acquiredIds.has(id)) { + id++; + } + + this.#acquiredIds.add(id); return id; } + + release(id) { + this.#acquiredIds.delete(id); + } } function createTestFileList(patterns, cwd) { @@ -538,94 +541,102 @@ function runTestFile(path, filesWatcher, opts) { debug('Assigned worker ID %d to test file: %s', workerId, path); } - if (watchMode) { - stdio.push('ipc'); - env.WATCH_REPORT_DEPENDENCIES = '1'; - } - if (opts.root.harness.shouldColorizeTestFiles) { - env.FORCE_COLOR = '1'; - } - - const child = spawn( - process.execPath, args, - { - __proto__: null, - signal: t.signal, - encoding: 'utf8', - env, - stdio, - cwd: opts.cwd, - }, - ); - if (watchMode) { - filesWatcher.runningProcesses.set(path, child); - filesWatcher.watcher.watchChildProcessModules(child, path); - } - - let err; + try { + if (watchMode) { + stdio.push('ipc'); + env.WATCH_REPORT_DEPENDENCIES = '1'; + } + if (opts.root.harness.shouldColorizeTestFiles) { + env.FORCE_COLOR = '1'; + } - child.on('error', (error) => { - err = error; - }); + const child = spawn( + process.execPath, args, + { + __proto__: null, + signal: t.signal, + encoding: 'utf8', + env, + stdio, + cwd: opts.cwd, + }, + ); + if (watchMode) { + filesWatcher.runningProcesses.set(path, child); + filesWatcher.watcher.watchChildProcessModules(child, path); + } - child.stdout.on('data', (data) => { - subtest.parseMessage(data); - }); + let err; - const rl = new Interface({ __proto__: null, input: child.stderr }); - rl.on('line', (line) => { - if (isInspectorMessage(line)) { - process.stderr.write(line + '\n'); - return; - } + child.on('error', (error) => { + err = error; + }); - // stderr cannot be treated as TAP, per the spec. However, we want to - // surface stderr lines to improve the DX. Inject each line into the - // test output as an unknown token as if it came from the TAP parser. - subtest.addToReport({ - __proto__: null, - type: 'test:stderr', - data: { __proto__: null, file: path, message: line + '\n' }, + child.stdout.on('data', (data) => { + subtest.parseMessage(data); }); - }); - const { 0: { 0: code, 1: signal } } = await SafePromiseAll([ - once(child, 'exit', { __proto__: null, signal: t.signal }), - finished(child.stdout, { __proto__: null, signal: t.signal }), - ]); - - // Close readline interface to prevent memory leak - rl.close(); - - if (watchMode) { - filesWatcher.runningProcesses.delete(path); - filesWatcher.runningSubtests.delete(path); - (async () => { - try { - await subTestEnded; - } finally { - if (filesWatcher.runningSubtests.size === 0) { - opts.root.reporter[kEmitMessage]('test:watch:drained'); - opts.root.postRun(); - } + const rl = new Interface({ __proto__: null, input: child.stderr }); + rl.on('line', (line) => { + if (isInspectorMessage(line)) { + process.stderr.write(line + '\n'); + return; } - })(); - } - if (code !== 0 || signal !== null) { - if (!err) { - const failureType = subtest.failedSubtests ? kSubtestsFailed : kTestCodeFailure; - err = ObjectAssign(new ERR_TEST_FAILURE('test failed', failureType), { + // stderr cannot be treated as TAP, per the spec. However, we want to + // surface stderr lines to improve the DX. Inject each line into the + // test output as an unknown token as if it came from the TAP parser. + subtest.addToReport({ __proto__: null, - exitCode: code, - signal: signal, - // The stack will not be useful since the failures came from tests - // in a child process. - stack: undefined, + type: 'test:stderr', + data: { __proto__: null, file: path, message: line + '\n' }, }); + }); + + const { 0: { 0: code, 1: signal } } = await SafePromiseAll([ + once(child, 'exit', { __proto__: null, signal: t.signal }), + finished(child.stdout, { __proto__: null, signal: t.signal }), + ]); + + // Close readline interface to prevent memory leak + rl.close(); + + if (watchMode) { + filesWatcher.runningProcesses.delete(path); + filesWatcher.runningSubtests.delete(path); + (async () => { + try { + await subTestEnded; + } finally { + if (filesWatcher.runningSubtests.size === 0) { + opts.root.reporter[kEmitMessage]('test:watch:drained'); + opts.root.postRun(); + } + } + })(); } - throw err; + if (code !== 0 || signal !== null) { + if (!err) { + const failureType = subtest.failedSubtests ? kSubtestsFailed : kTestCodeFailure; + err = ObjectAssign(new ERR_TEST_FAILURE('test failed', failureType), { + __proto__: null, + exitCode: code, + signal: signal, + // The stack will not be useful since the failures came from tests + // in a child process. + stack: undefined, + }); + } + + throw err; + } + } finally { + // Every exit path must return the ID, including abort and spawn failure. + if (opts.workerIdPool && workerId !== undefined) { + opts.workerIdPool.release(workerId); + debug('Released worker ID %d from test file: %s', workerId, path); + } } }); const subTestEnded = subtest.start(); @@ -1011,23 +1022,10 @@ function run(options = kEmptyObject) { let filesWatcher; let runFiles; - // Create worker ID pool for concurrent test execution. - // Use concurrency from globalOptions which has been processed by parseCommandLine(). - const effectiveConcurrency = globalOptions.concurrency ?? concurrency; - let maxConcurrency = 1; - if (effectiveConcurrency === true) { - maxConcurrency = MathMax(availableParallelism() - 1, 1); - } else if (typeof effectiveConcurrency === 'number') { - maxConcurrency = effectiveConcurrency; - } - const workerIdPool = new WorkerIdPool(maxConcurrency); - debug( - 'Created worker ID pool with max concurrency: %d, ' + - 'effectiveConcurrency: %s, testFiles: %d', - maxConcurrency, - effectiveConcurrency, - testFiles.length, - ); + // The pool tracks the IDs actually in use, so they stay exclusive and never + // exceed the number of files running concurrently. + const workerIdPool = new WorkerIdPool(); + debug('Created worker ID pool, testFiles: %d', testFiles.length); const opts = { __proto__: null, diff --git a/test/parallel/test-runner-worker-id.js b/test/parallel/test-runner-worker-id.js index 9366cfb4fd8a..0bb21dadcf2a 100644 --- a/test/parallel/test-runner-worker-id.js +++ b/test/parallel/test-runner-worker-id.js @@ -1,8 +1,10 @@ 'use strict'; -require('../common'); +const common = require('../common'); +const tmpdir = require('../common/tmpdir'); const fixtures = require('../common/fixtures'); const assert = require('node:assert'); const { spawnSync } = require('node:child_process'); +const { readFileSync, writeFileSync } = require('node:fs'); const { test } = require('node:test'); test('NODE_TEST_WORKER_ID is set for concurrent test files', async () => { @@ -87,8 +89,6 @@ test('context.workerId matches NODE_TEST_WORKER_ID', async () => { }); test('worker IDs are reused when more tests than concurrency', async () => { - const tmpdir = require('../common/tmpdir'); - const { writeFileSync } = require('node:fs'); tmpdir.refresh(); // Create 9 separate test files dynamically @@ -119,7 +119,6 @@ test('track worker ${i}', () => { assert.strictEqual(result.status, 0, `Test failed: ${result.stderr.toString()}`); // Read and analyze worker IDs used - const { readFileSync } = require('node:fs'); const workerIds = readFileSync(usageFile, 'utf8').trim().split('\n'); // Count occurrences of each worker ID @@ -139,3 +138,106 @@ test('track worker ${i}', () => { assert.strictEqual(count, 3, `Worker ID ${id} should be used 3 times, got ${count}`); }); }); + +// Generates a test file that appends ` ` lines to a +// shared log, which is replayed below to find IDs held by two live files at +// once. +function testFile(name, { block = false, release = false } = {}) { + let body = ''; + + if (release) { + body += "writeFileSync(process.env.MARKER_FILE, '');\n"; + } + if (block) { + const deadline = common.platformTimeout(10_000); + body += `const buf = new Int32Array(new SharedArrayBuffer(4)); + const deadline = Date.now() + ${deadline}; + while (!existsSync(process.env.MARKER_FILE) && Date.now() < deadline) { + Atomics.wait(buf, 0, 0, 20); + }\n`; + } + + return ` +import { test } from 'node:test'; +import { appendFileSync, writeFileSync, existsSync } from 'node:fs'; + +const log = (kind) => appendFileSync(process.env.WORKER_LOG, + kind + ' ${name} ' + process.env.NODE_TEST_WORKER_ID + '\\n'); + +test('${name}', () => { + log('start'); + ${body} + log('end'); +}); +`; +} + +// Replays the log and reports every ID that was held by two files at once. +function findConflicts(events) { + const live = new Map(); + const conflicts = []; + + for (const event of events) { + const [kind, name, id] = event.split(' '); + if (kind === 'start') { + if (live.has(id)) { + conflicts.push(`worker ID ${id} held by both '${live.get(id)}' and '${name}'`); + } + live.set(id, name); + } else { + live.delete(id); + } + } + + return conflicts; +} + +test('worker IDs are exclusive to concurrently running test files', () => { + tmpdir.refresh(); + + // `slow` pins one ID for the whole run while the remaining files churn + // through the other slots, so IDs are released and reacquired several times + // with one of them permanently taken. + const sources = [['slow', { block: true }]]; + for (let i = 1; i <= 4; i++) { + sources.push([`file-${i}`]); + } + sources.push(['last', { release: true }]); + + const concurrency = 3; + const logFile = tmpdir.resolve('worker-id-log.txt'); + const markerFile = tmpdir.resolve('worker-id-release.marker'); + writeFileSync(logFile, ''); + + const files = sources.map(([name, options], index) => { + const file = tmpdir.resolve(`worker-id-${index}-${name}.mjs`); + writeFileSync(file, testFile(name, options)); + return file; + }); + + const result = spawnSync( + process.execPath, + ['--test', `--test-concurrency=${concurrency}`, ...files], + { env: { ...process.env, WORKER_LOG: logFile, MARKER_FILE: markerFile } }, + ); + assert.strictEqual(result.status, 0, `Runner failed: ${result.stderr.toString()}`); + + const events = readFileSync(logFile, 'utf8').trim().split('\n'); + assert.strictEqual(events.length, sources.length * 2, + `Unexpected event log:\n${events.join('\n')}`); + + // The blocking file must still be running when the last file starts + const eventsAt = (line) => events.findIndex((event) => event.startsWith(line)); + assert.ok(eventsAt('start last') < eventsAt('end slow'), + `Files did not overlap:\n${events.join('\n')}`); + + const conflicts = findConflicts(events); + assert.deepStrictEqual( + conflicts, [], `Worker IDs were not exclusive:\n${events.join('\n')}\n\n${conflicts.join('\n')}`); + + for (const event of events) { + const id = Number(event.split(' ')[2]); + assert.ok(Number.isInteger(id) && id >= 1 && id <= concurrency, + `Worker ID outside 1..${concurrency}:\n${events.join('\n')}`); + } +});