Skip to content

Commit 5416f6a

Browse files
committed
test_runner: do not reuse a worker ID held by a running file
Worker IDs were handed out round-robin and never released, so a file that started after another finished could get an ID still held by a live process. Track the IDs in use, hand out the lowest free one, and release it in a finally block once the child process exits. Refs: #61394 Signed-off-by: Vasiliy Serpokryl <vasiliy.serpokryl@mail.ru>
1 parent 5ed55ba commit 5416f6a

3 files changed

Lines changed: 209 additions & 107 deletions

File tree

‎doc/api/test.md‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4575,7 +4575,9 @@ The unique identifier of the worker running the current test file. This value is
45754575
derived from the `NODE_TEST_WORKER_ID` environment variable. When running tests
45764576
with `--test-isolation=process` (the default), each test file runs in a separate
45774577
child process and is assigned a worker ID from 1 to N, where N is the number of
4578-
concurrent workers. When running with `--test-isolation=none`, all tests run in
4578+
concurrent workers. A worker ID is never shared by two test files running at the
4579+
same time. Once a test file finishes, its worker ID is reused by the next test
4580+
file that starts. When running with `--test-isolation=none`, all tests run in
45794581
the same process and the worker ID is always 1. This value is `undefined` when
45804582
not running in a test context.
45814583

‎lib/internal/test_runner/runner.js‎

Lines changed: 100 additions & 102 deletions
Original file line numberDiff line numberDiff line change
@@ -15,7 +15,6 @@ const {
1515
ArrayPrototypeSlice,
1616
ArrayPrototypeSome,
1717
ArrayPrototypeSort,
18-
MathMax,
1918
ObjectAssign,
2019
PromisePrototypeThen,
2120
PromiseWithResolvers,
@@ -38,7 +37,6 @@ const {
3837
const { spawn } = require('child_process');
3938
const { statSync } = require('fs');
4039
const { finished } = require('internal/streams/end-of-stream');
41-
const { availableParallelism } = require('os');
4240
const { resolve, sep, isAbsolute } = require('path');
4341
const { DefaultDeserializer, DefaultSerializer } = require('v8');
4442
const { getOptionValue, getOptionsAsFlagsFromBinding } = require('internal/options');
@@ -139,17 +137,22 @@ let kResistStopPropagation;
139137

140138
// Worker ID pool management for concurrent test execution
141139
class WorkerIdPool {
142-
#nextId = 0;
143-
#maxConcurrency;
144-
145-
constructor(maxConcurrency) {
146-
this.#maxConcurrency = maxConcurrency;
147-
}
140+
#acquiredIds = new SafeSet();
148141

149142
acquire() {
150-
const id = (this.#nextId++ % this.#maxConcurrency) + 1;
143+
let id = 1;
144+
145+
while (this.#acquiredIds.has(id)) {
146+
id++;
147+
}
148+
149+
this.#acquiredIds.add(id);
151150
return id;
152151
}
152+
153+
release(id) {
154+
this.#acquiredIds.delete(id);
155+
}
153156
}
154157

155158
function createTestFileList(patterns, cwd) {
@@ -538,94 +541,102 @@ function runTestFile(path, filesWatcher, opts) {
538541
debug('Assigned worker ID %d to test file: %s', workerId, path);
539542
}
540543

541-
if (watchMode) {
542-
stdio.push('ipc');
543-
env.WATCH_REPORT_DEPENDENCIES = '1';
544-
}
545-
if (opts.root.harness.shouldColorizeTestFiles) {
546-
env.FORCE_COLOR = '1';
547-
}
548-
549-
const child = spawn(
550-
process.execPath, args,
551-
{
552-
__proto__: null,
553-
signal: t.signal,
554-
encoding: 'utf8',
555-
env,
556-
stdio,
557-
cwd: opts.cwd,
558-
},
559-
);
560-
if (watchMode) {
561-
filesWatcher.runningProcesses.set(path, child);
562-
filesWatcher.watcher.watchChildProcessModules(child, path);
563-
}
564-
565-
let err;
544+
try {
545+
if (watchMode) {
546+
stdio.push('ipc');
547+
env.WATCH_REPORT_DEPENDENCIES = '1';
548+
}
549+
if (opts.root.harness.shouldColorizeTestFiles) {
550+
env.FORCE_COLOR = '1';
551+
}
566552

567-
child.on('error', (error) => {
568-
err = error;
569-
});
553+
const child = spawn(
554+
process.execPath, args,
555+
{
556+
__proto__: null,
557+
signal: t.signal,
558+
encoding: 'utf8',
559+
env,
560+
stdio,
561+
cwd: opts.cwd,
562+
},
563+
);
564+
if (watchMode) {
565+
filesWatcher.runningProcesses.set(path, child);
566+
filesWatcher.watcher.watchChildProcessModules(child, path);
567+
}
570568

571-
child.stdout.on('data', (data) => {
572-
subtest.parseMessage(data);
573-
});
569+
let err;
574570

575-
const rl = new Interface({ __proto__: null, input: child.stderr });
576-
rl.on('line', (line) => {
577-
if (isInspectorMessage(line)) {
578-
process.stderr.write(line + '\n');
579-
return;
580-
}
571+
child.on('error', (error) => {
572+
err = error;
573+
});
581574

582-
// stderr cannot be treated as TAP, per the spec. However, we want to
583-
// surface stderr lines to improve the DX. Inject each line into the
584-
// test output as an unknown token as if it came from the TAP parser.
585-
subtest.addToReport({
586-
__proto__: null,
587-
type: 'test:stderr',
588-
data: { __proto__: null, file: path, message: line + '\n' },
575+
child.stdout.on('data', (data) => {
576+
subtest.parseMessage(data);
589577
});
590-
});
591578

592-
const { 0: { 0: code, 1: signal } } = await SafePromiseAll([
593-
once(child, 'exit', { __proto__: null, signal: t.signal }),
594-
finished(child.stdout, { __proto__: null, signal: t.signal }),
595-
]);
596-
597-
// Close readline interface to prevent memory leak
598-
rl.close();
599-
600-
if (watchMode) {
601-
filesWatcher.runningProcesses.delete(path);
602-
filesWatcher.runningSubtests.delete(path);
603-
(async () => {
604-
try {
605-
await subTestEnded;
606-
} finally {
607-
if (filesWatcher.runningSubtests.size === 0) {
608-
opts.root.reporter[kEmitMessage]('test:watch:drained');
609-
opts.root.postRun();
610-
}
579+
const rl = new Interface({ __proto__: null, input: child.stderr });
580+
rl.on('line', (line) => {
581+
if (isInspectorMessage(line)) {
582+
process.stderr.write(line + '\n');
583+
return;
611584
}
612-
})();
613-
}
614585

615-
if (code !== 0 || signal !== null) {
616-
if (!err) {
617-
const failureType = subtest.failedSubtests ? kSubtestsFailed : kTestCodeFailure;
618-
err = ObjectAssign(new ERR_TEST_FAILURE('test failed', failureType), {
586+
// stderr cannot be treated as TAP, per the spec. However, we want to
587+
// surface stderr lines to improve the DX. Inject each line into the
588+
// test output as an unknown token as if it came from the TAP parser.
589+
subtest.addToReport({
619590
__proto__: null,
620-
exitCode: code,
621-
signal: signal,
622-
// The stack will not be useful since the failures came from tests
623-
// in a child process.
624-
stack: undefined,
591+
type: 'test:stderr',
592+
data: { __proto__: null, file: path, message: line + '\n' },
625593
});
594+
});
595+
596+
const { 0: { 0: code, 1: signal } } = await SafePromiseAll([
597+
once(child, 'exit', { __proto__: null, signal: t.signal }),
598+
finished(child.stdout, { __proto__: null, signal: t.signal }),
599+
]);
600+
601+
// Close readline interface to prevent memory leak
602+
rl.close();
603+
604+
if (watchMode) {
605+
filesWatcher.runningProcesses.delete(path);
606+
filesWatcher.runningSubtests.delete(path);
607+
(async () => {
608+
try {
609+
await subTestEnded;
610+
} finally {
611+
if (filesWatcher.runningSubtests.size === 0) {
612+
opts.root.reporter[kEmitMessage]('test:watch:drained');
613+
opts.root.postRun();
614+
}
615+
}
616+
})();
626617
}
627618

628-
throw err;
619+
if (code !== 0 || signal !== null) {
620+
if (!err) {
621+
const failureType = subtest.failedSubtests ? kSubtestsFailed : kTestCodeFailure;
622+
err = ObjectAssign(new ERR_TEST_FAILURE('test failed', failureType), {
623+
__proto__: null,
624+
exitCode: code,
625+
signal: signal,
626+
// The stack will not be useful since the failures came from tests
627+
// in a child process.
628+
stack: undefined,
629+
});
630+
}
631+
632+
throw err;
633+
}
634+
} finally {
635+
// Every exit path must return the ID, including abort and spawn failure.
636+
if (opts.workerIdPool && workerId !== undefined) {
637+
opts.workerIdPool.release(workerId);
638+
debug('Released worker ID %d from test file: %s', workerId, path);
639+
}
629640
}
630641
});
631642
const subTestEnded = subtest.start();
@@ -1011,23 +1022,10 @@ function run(options = kEmptyObject) {
10111022
let filesWatcher;
10121023
let runFiles;
10131024

1014-
// Create worker ID pool for concurrent test execution.
1015-
// Use concurrency from globalOptions which has been processed by parseCommandLine().
1016-
const effectiveConcurrency = globalOptions.concurrency ?? concurrency;
1017-
let maxConcurrency = 1;
1018-
if (effectiveConcurrency === true) {
1019-
maxConcurrency = MathMax(availableParallelism() - 1, 1);
1020-
} else if (typeof effectiveConcurrency === 'number') {
1021-
maxConcurrency = effectiveConcurrency;
1022-
}
1023-
const workerIdPool = new WorkerIdPool(maxConcurrency);
1024-
debug(
1025-
'Created worker ID pool with max concurrency: %d, ' +
1026-
'effectiveConcurrency: %s, testFiles: %d',
1027-
maxConcurrency,
1028-
effectiveConcurrency,
1029-
testFiles.length,
1030-
);
1025+
// The pool tracks the IDs actually in use, so they stay exclusive and never
1026+
// exceed the number of files running concurrently.
1027+
const workerIdPool = new WorkerIdPool();
1028+
debug('Created worker ID pool, testFiles: %d', testFiles.length);
10311029

10321030
const opts = {
10331031
__proto__: null,

‎test/parallel/test-runner-worker-id.js‎

Lines changed: 106 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1,8 +1,10 @@
11
'use strict';
2-
require('../common');
2+
const common = require('../common');
3+
const tmpdir = require('../common/tmpdir');
34
const fixtures = require('../common/fixtures');
45
const assert = require('node:assert');
56
const { spawnSync } = require('node:child_process');
7+
const { readFileSync, writeFileSync } = require('node:fs');
68
const { test } = require('node:test');
79

810
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 () => {
8789
});
8890

8991
test('worker IDs are reused when more tests than concurrency', async () => {
90-
const tmpdir = require('../common/tmpdir');
91-
const { writeFileSync } = require('node:fs');
9292
tmpdir.refresh();
9393

9494
// Create 9 separate test files dynamically
@@ -119,7 +119,6 @@ test('track worker ${i}', () => {
119119
assert.strictEqual(result.status, 0, `Test failed: ${result.stderr.toString()}`);
120120

121121
// Read and analyze worker IDs used
122-
const { readFileSync } = require('node:fs');
123122
const workerIds = readFileSync(usageFile, 'utf8').trim().split('\n');
124123

125124
// Count occurrences of each worker ID
@@ -139,3 +138,106 @@ test('track worker ${i}', () => {
139138
assert.strictEqual(count, 3, `Worker ID ${id} should be used 3 times, got ${count}`);
140139
});
141140
});
141+
142+
// Generates a test file that appends `<kind> <name> <workerId>` lines to a
143+
// shared log, which is replayed below to find IDs held by two live files at
144+
// once.
145+
function testFile(name, { block = false, release = false } = {}) {
146+
let body = '';
147+
148+
if (release) {
149+
body += "writeFileSync(process.env.MARKER_FILE, '');\n";
150+
}
151+
if (block) {
152+
const deadline = common.platformTimeout(10_000);
153+
body += `const buf = new Int32Array(new SharedArrayBuffer(4));
154+
const deadline = Date.now() + ${deadline};
155+
while (!existsSync(process.env.MARKER_FILE) && Date.now() < deadline) {
156+
Atomics.wait(buf, 0, 0, 20);
157+
}\n`;
158+
}
159+
160+
return `
161+
import { test } from 'node:test';
162+
import { appendFileSync, writeFileSync, existsSync } from 'node:fs';
163+
164+
const log = (kind) => appendFileSync(process.env.WORKER_LOG,
165+
kind + ' ${name} ' + process.env.NODE_TEST_WORKER_ID + '\\n');
166+
167+
test('${name}', () => {
168+
log('start');
169+
${body}
170+
log('end');
171+
});
172+
`;
173+
}
174+
175+
// Replays the log and reports every ID that was held by two files at once.
176+
function findConflicts(events) {
177+
const live = new Map();
178+
const conflicts = [];
179+
180+
for (const event of events) {
181+
const [kind, name, id] = event.split(' ');
182+
if (kind === 'start') {
183+
if (live.has(id)) {
184+
conflicts.push(`worker ID ${id} held by both '${live.get(id)}' and '${name}'`);
185+
}
186+
live.set(id, name);
187+
} else {
188+
live.delete(id);
189+
}
190+
}
191+
192+
return conflicts;
193+
}
194+
195+
test('worker IDs are exclusive to concurrently running test files', () => {
196+
tmpdir.refresh();
197+
198+
// `slow` pins one ID for the whole run while the remaining files churn
199+
// through the other slots, so IDs are released and reacquired several times
200+
// with one of them permanently taken.
201+
const sources = [['slow', { block: true }]];
202+
for (let i = 1; i <= 4; i++) {
203+
sources.push([`file-${i}`]);
204+
}
205+
sources.push(['last', { release: true }]);
206+
207+
const concurrency = 3;
208+
const logFile = tmpdir.resolve('worker-id-log.txt');
209+
const markerFile = tmpdir.resolve('worker-id-release.marker');
210+
writeFileSync(logFile, '');
211+
212+
const files = sources.map(([name, options], index) => {
213+
const file = tmpdir.resolve(`worker-id-${index}-${name}.mjs`);
214+
writeFileSync(file, testFile(name, options));
215+
return file;
216+
});
217+
218+
const result = spawnSync(
219+
process.execPath,
220+
['--test', `--test-concurrency=${concurrency}`, ...files],
221+
{ env: { ...process.env, WORKER_LOG: logFile, MARKER_FILE: markerFile } },
222+
);
223+
assert.strictEqual(result.status, 0, `Runner failed: ${result.stderr.toString()}`);
224+
225+
const events = readFileSync(logFile, 'utf8').trim().split('\n');
226+
assert.strictEqual(events.length, sources.length * 2,
227+
`Unexpected event log:\n${events.join('\n')}`);
228+
229+
// The blocking file must still be running when the last file starts
230+
const eventsAt = (line) => events.findIndex((event) => event.startsWith(line));
231+
assert.ok(eventsAt('start last') < eventsAt('end slow'),
232+
`Files did not overlap:\n${events.join('\n')}`);
233+
234+
const conflicts = findConflicts(events);
235+
assert.deepStrictEqual(
236+
conflicts, [], `Worker IDs were not exclusive:\n${events.join('\n')}\n\n${conflicts.join('\n')}`);
237+
238+
for (const event of events) {
239+
const id = Number(event.split(' ')[2]);
240+
assert.ok(Number.isInteger(id) && id >= 1 && id <= concurrency,
241+
`Worker ID outside 1..${concurrency}:\n${events.join('\n')}`);
242+
}
243+
});

0 commit comments

Comments
 (0)