From 4362afbb6f5287d2e5414e5b3b413b96870bab7f Mon Sep 17 00:00:00 2001 From: Filip Skokan Date: Thu, 1 Oct 2026 14:10:35 +0200 Subject: [PATCH 01/10] test: run suites in parallel by default Make serial execution an explicit suite choice. Run sequential tests with the normal worker count while allowing only one test per filename subsystem, retaining ownership through retries and releasing it on shutdown. Finish parallel tests before serial suites and subsystem-concurrent tests. Preserve the existing scheduling safeguards for abort, wasm, WPT and SEA. Signed-off-by: Filip Skokan Assisted-by: Codex --- test/README.md | 17 +- test/abort/testcfg.py | 1 + test/addons/testcfg.py | 2 +- test/async-hooks/testcfg.py | 2 +- test/benchmark/testcfg.py | 3 +- test/client-proxy/testcfg.py | 2 +- test/doctool/testcfg.py | 2 +- test/embedding/testcfg.py | 2 +- test/es-module/testcfg.py | 2 +- test/ffi/testcfg.py | 2 +- test/internet/testcfg.py | 3 +- test/js-native-api/testcfg.py | 2 +- test/known_issues/testcfg.py | 2 +- test/module-hooks/testcfg.py | 2 +- test/node-api/testcfg.py | 2 +- test/parallel/testcfg.py | 2 +- test/pummel/testcfg.py | 2 +- test/report/testcfg.py | 2 +- test/sea/testcfg.py | 7 +- test/sequential/testcfg.py | 2 +- test/system-ca/test.cfg.py | 2 +- test/test-runner/testcfg.py | 2 +- test/test426/testcfg.py | 2 +- test/testpy/__init__.py | 22 +- test/tools/test_test_configurations.py | 116 +++++++ test/tools/test_test_runner.py | 444 +++++++++++++++++++++++++ test/trace_events/testcfg.py | 2 +- test/v8-updates/testcfg.py | 2 +- test/wasi/testcfg.py | 2 +- test/wasm-allocation/testcfg.py | 3 +- test/wpt/testcfg.py | 3 +- tools/test.py | 208 ++++++++---- 32 files changed, 770 insertions(+), 99 deletions(-) create mode 100644 test/tools/test_test_configurations.py create mode 100644 test/tools/test_test_runner.py diff --git a/test/README.md b/test/README.md index 02359ca68265..37d54c544c98 100644 --- a/test/README.md +++ b/test/README.md @@ -33,11 +33,26 @@ For the tests to run on Windows, be sure to clone Node.js source code with the | `parallel` | Yes | Various tests that are able to be run in parallel. | | `pseudo-tty` | Yes | Tests that require stdin/stdout/stderr to be a TTY. | | `pummel` | No | Various tests for various modules / system functionality operating under load. | -| `sequential` | Yes | Various tests that must not run in parallel. | +| `sequential` | Yes | Tests that run sequentially within each filename subsystem. | | `testpy` | _N/A_ | Test configuration utility used by various test suites. | | `tick-processor` | No | Tests for the V8 tick processor integration.[^4] | | `v8-updates` | No | Tests for V8 performance integration. | +Tests run in parallel by default. Suite configurations opt out with +`SerialTestConfiguration`, or `SerialAddonTestConfiguration` for addon layouts. +The `addons`, `benchmark`, `internet`, `js-native-api`, `known_issues`, `node-api`, +and `pummel` suites explicitly run serially. The `abort` and `wasm-allocation` +suites also remain serial; WPT's group settings and SEA's disk-space guard retain +their existing scheduling. + +The test runner finishes parallel tests first, followed by serial suites, then +`sequential` tests. In `sequential`, different subsystems can run concurrently +up to the worker count selected with `-j`. The subsystem is the first filename +component after `test-`, so `test-net-server-bind.js` and +`test-net-connect-econnrefused.js` cannot overlap, while a `test-fs-*` test can run +alongside them. This scheduling does not isolate resources shared across +subsystems, such as fixed ports. + [^1]: [Documentation](../test/common/README.md) [^2]: Tests for networking related modules may also be present in other directories, but those tests do diff --git a/test/abort/testcfg.py b/test/abort/testcfg.py index e509d0453c40..40b1ab429175 100644 --- a/test/abort/testcfg.py +++ b/test/abort/testcfg.py @@ -3,4 +3,5 @@ import testpy def GetConfiguration(context, root): + # TODO: Replace preexec_fn core suppression before parallelizing this suite. return testpy.AbortTestConfiguration(context, root, 'abort') diff --git a/test/addons/testcfg.py b/test/addons/testcfg.py index 6c61081fbb0b..15924019b700 100644 --- a/test/addons/testcfg.py +++ b/test/addons/testcfg.py @@ -3,4 +3,4 @@ import testpy def GetConfiguration(context, root): - return testpy.AddonTestConfiguration(context, root, 'addons') + return testpy.SerialAddonTestConfiguration(context, root, 'addons') diff --git a/test/async-hooks/testcfg.py b/test/async-hooks/testcfg.py index 79e8c6834eb3..be4158c46518 100644 --- a/test/async-hooks/testcfg.py +++ b/test/async-hooks/testcfg.py @@ -3,4 +3,4 @@ import testpy def GetConfiguration(context, root): - return testpy.ParallelTestConfiguration(context, root, 'async-hooks') + return testpy.SimpleTestConfiguration(context, root, 'async-hooks') diff --git a/test/benchmark/testcfg.py b/test/benchmark/testcfg.py index 2c2929f610b8..cf364de83f0b 100644 --- a/test/benchmark/testcfg.py +++ b/test/benchmark/testcfg.py @@ -3,4 +3,5 @@ import testpy def GetConfiguration(context, root): - return testpy.SimpleTestConfiguration(context, root, 'benchmark') + # TODO: Isolate fixed TCP/UDP ports before parallelizing benchmark tests. + return testpy.SerialTestConfiguration(context, root, 'benchmark') diff --git a/test/client-proxy/testcfg.py b/test/client-proxy/testcfg.py index 055549e7a04c..7daedd2319a4 100644 --- a/test/client-proxy/testcfg.py +++ b/test/client-proxy/testcfg.py @@ -3,4 +3,4 @@ import testpy def GetConfiguration(context, root): - return testpy.ParallelTestConfiguration(context, root, 'client-proxy') + return testpy.SimpleTestConfiguration(context, root, 'client-proxy') diff --git a/test/doctool/testcfg.py b/test/doctool/testcfg.py index 33a274a43d64..5778d2f0c5ba 100644 --- a/test/doctool/testcfg.py +++ b/test/doctool/testcfg.py @@ -4,4 +4,4 @@ import testpy def GetConfiguration(context, root): - return testpy.ParallelTestConfiguration(context, root, 'doctool') + return testpy.SimpleTestConfiguration(context, root, 'doctool') diff --git a/test/embedding/testcfg.py b/test/embedding/testcfg.py index f63bc393a6a8..a4b90f490c67 100644 --- a/test/embedding/testcfg.py +++ b/test/embedding/testcfg.py @@ -3,4 +3,4 @@ import testpy def GetConfiguration(context, root): - return testpy.ParallelTestConfiguration(context, root, 'embedding') + return testpy.SimpleTestConfiguration(context, root, 'embedding') diff --git a/test/es-module/testcfg.py b/test/es-module/testcfg.py index 83ce46ee6672..0d8dfeed463e 100644 --- a/test/es-module/testcfg.py +++ b/test/es-module/testcfg.py @@ -3,4 +3,4 @@ import testpy def GetConfiguration(context, root): - return testpy.ParallelTestConfiguration(context, root, 'es-module') + return testpy.SimpleTestConfiguration(context, root, 'es-module') diff --git a/test/ffi/testcfg.py b/test/ffi/testcfg.py index 3d51dc261db9..75ac86fe91e4 100644 --- a/test/ffi/testcfg.py +++ b/test/ffi/testcfg.py @@ -3,4 +3,4 @@ import testpy def GetConfiguration(context, root): - return testpy.ParallelTestConfiguration(context, root, 'ffi') + return testpy.SimpleTestConfiguration(context, root, 'ffi') diff --git a/test/internet/testcfg.py b/test/internet/testcfg.py index 73e70e340000..7c68a8573976 100644 --- a/test/internet/testcfg.py +++ b/test/internet/testcfg.py @@ -3,4 +3,5 @@ import testpy def GetConfiguration(context, root): - return testpy.SimpleTestConfiguration(context, root, 'internet') + # TODO: Isolate shared listening ports before allowing concurrent tests. + return testpy.SerialTestConfiguration(context, root, 'internet') diff --git a/test/js-native-api/testcfg.py b/test/js-native-api/testcfg.py index 4e5d67709a87..ca189cf6081d 100644 --- a/test/js-native-api/testcfg.py +++ b/test/js-native-api/testcfg.py @@ -3,4 +3,4 @@ import testpy def GetConfiguration(context, root): - return testpy.AddonTestConfiguration(context, root, 'js-native-api') + return testpy.SerialAddonTestConfiguration(context, root, 'js-native-api') diff --git a/test/known_issues/testcfg.py b/test/known_issues/testcfg.py index d8d2ba22b570..6a66b40c631c 100644 --- a/test/known_issues/testcfg.py +++ b/test/known_issues/testcfg.py @@ -6,4 +6,4 @@ def GetConfiguration(context, root): myContext = copy.copy(context) myContext.expect_fail = 1 - return testpy.SimpleTestConfiguration(myContext, root, 'known_issues') + return testpy.SerialTestConfiguration(myContext, root, 'known_issues') diff --git a/test/module-hooks/testcfg.py b/test/module-hooks/testcfg.py index f904b1e9170f..39795b0ff6cd 100644 --- a/test/module-hooks/testcfg.py +++ b/test/module-hooks/testcfg.py @@ -3,4 +3,4 @@ import testpy def GetConfiguration(context, root): - return testpy.ParallelTestConfiguration(context, root, 'module-hooks') + return testpy.SimpleTestConfiguration(context, root, 'module-hooks') diff --git a/test/node-api/testcfg.py b/test/node-api/testcfg.py index 3453f3bc3738..f720336cf297 100644 --- a/test/node-api/testcfg.py +++ b/test/node-api/testcfg.py @@ -3,4 +3,4 @@ import testpy def GetConfiguration(context, root): - return testpy.AddonTestConfiguration(context, root, 'node-api') + return testpy.SerialAddonTestConfiguration(context, root, 'node-api') diff --git a/test/parallel/testcfg.py b/test/parallel/testcfg.py index 8b610b50b1d7..22283d9e2cf5 100644 --- a/test/parallel/testcfg.py +++ b/test/parallel/testcfg.py @@ -3,4 +3,4 @@ import testpy def GetConfiguration(context, root): - return testpy.ParallelTestConfiguration(context, root, 'parallel') + return testpy.SimpleTestConfiguration(context, root, 'parallel') diff --git a/test/pummel/testcfg.py b/test/pummel/testcfg.py index c91fe39dce27..23fb92b3a387 100644 --- a/test/pummel/testcfg.py +++ b/test/pummel/testcfg.py @@ -3,4 +3,4 @@ import testpy def GetConfiguration(context, root): - return testpy.SimpleTestConfiguration(context, root, 'pummel') + return testpy.SerialTestConfiguration(context, root, 'pummel') diff --git a/test/report/testcfg.py b/test/report/testcfg.py index c06b75ce5c7d..8c1f7f24b18f 100644 --- a/test/report/testcfg.py +++ b/test/report/testcfg.py @@ -3,4 +3,4 @@ import testpy def GetConfiguration(context, root): - return testpy.ParallelTestConfiguration(context, root, 'report') + return testpy.SimpleTestConfiguration(context, root, 'report') diff --git a/test/sea/testcfg.py b/test/sea/testcfg.py index 8f3acefe4ece..bc08b157add5 100644 --- a/test/sea/testcfg.py +++ b/test/sea/testcfg.py @@ -15,14 +15,15 @@ def GetConfiguration(context, root): vm = context.GetVm('none', preferred_mode) if not os.path.isfile(vm): - return testpy.SimpleTestConfiguration(context, root, 'sea') + return testpy.SerialTestConfiguration(context, root, 'sea') # Get the size of the executable to decide whether we can run tests in parallel. + # TODO: Use the requested worker count when evaluating this disk-space guard. executable_size = os.path.getsize(vm) num_cpus = multiprocessing.cpu_count() remaining_disk_space = shutil.disk_usage('.').free # Give it a bit of leeway by multiplying by 3. if (executable_size * num_cpus * 3 > remaining_disk_space): - return testpy.SimpleTestConfiguration(context, root, 'sea') + return testpy.SerialTestConfiguration(context, root, 'sea') - return testpy.ParallelTestConfiguration(context, root, 'sea') + return testpy.SimpleTestConfiguration(context, root, 'sea') diff --git a/test/sequential/testcfg.py b/test/sequential/testcfg.py index b1fce1810017..440e1ee82487 100644 --- a/test/sequential/testcfg.py +++ b/test/sequential/testcfg.py @@ -3,4 +3,4 @@ import testpy def GetConfiguration(context, root): - return testpy.SimpleTestConfiguration(context, root, 'sequential') + return testpy.SerialTestConfiguration(context, root, 'sequential') diff --git a/test/system-ca/test.cfg.py b/test/system-ca/test.cfg.py index 5b4d3fd1ab6e..d348149de964 100644 --- a/test/system-ca/test.cfg.py +++ b/test/system-ca/test.cfg.py @@ -3,4 +3,4 @@ import testpy def GetConfiguration(context, root): - return testpy.ParallelTestConfiguration(context, root, 'system-ca') + return testpy.SimpleTestConfiguration(context, root, 'system-ca') diff --git a/test/test-runner/testcfg.py b/test/test-runner/testcfg.py index fa53433d0dd1..dc5a484b6894 100644 --- a/test/test-runner/testcfg.py +++ b/test/test-runner/testcfg.py @@ -3,4 +3,4 @@ import testpy def GetConfiguration(context, root): - return testpy.ParallelTestConfiguration(context, root, 'test-runner') + return testpy.SimpleTestConfiguration(context, root, 'test-runner') diff --git a/test/test426/testcfg.py b/test/test426/testcfg.py index 235311f0dddd..18bfa7ea5b86 100644 --- a/test/test426/testcfg.py +++ b/test/test426/testcfg.py @@ -3,4 +3,4 @@ import testpy def GetConfiguration(context, root): - return testpy.ParallelTestConfiguration(context, root, 'test426') + return testpy.SimpleTestConfiguration(context, root, 'test426') diff --git a/test/testpy/__init__.py b/test/testpy/__init__.py index 468f94fe5c0a..816b6437490b 100644 --- a/test/testpy/__init__.py +++ b/test/testpy/__init__.py @@ -44,6 +44,7 @@ def __init__(self, path, file, arch, mode, context, config, additional=None): super(SimpleTestCase, self).__init__(context, path, arch, mode) self.file = file self.config = config + self.parallel = config.parallel self.arch = arch self.mode = mode if additional is not None: @@ -128,6 +129,8 @@ def GetSource(self): class SimpleTestConfiguration(test.TestConfiguration): + parallel = True + def __init__(self, context, root, section, additional=None): super(SimpleTestConfiguration, self).__init__(context, root, section) if additional is not None: @@ -152,17 +155,8 @@ def ListTests(self, current_path, path, arch, mode): def GetBuildRequirements(self): return ['sample', 'sample=shell'] -class ParallelTestConfiguration(SimpleTestConfiguration): - def __init__(self, context, root, section, additional=None): - super(ParallelTestConfiguration, self).__init__(context, root, section, - additional) - - def ListTests(self, current_path, path, arch, mode): - result = super(ParallelTestConfiguration, self).ListTests( - current_path, path, arch, mode) - for tst in result: - tst.parallel = True - return result +class SerialTestConfiguration(SimpleTestConfiguration): + parallel = False class AddonTestConfiguration(SimpleTestConfiguration): def __init__(self, context, root, section, additional=None): @@ -188,7 +182,11 @@ def ListTests(self, current_path, path, arch, mode): SimpleTestCase(tst, file_path, arch, mode, self.context, self, self.additional_flags)) return result -class AbortTestConfiguration(SimpleTestConfiguration): +class SerialAddonTestConfiguration(AddonTestConfiguration): + parallel = False + + +class AbortTestConfiguration(SerialTestConfiguration): def __init__(self, context, root, section, additional=None): super(AbortTestConfiguration, self).__init__(context, root, section, additional) diff --git a/test/tools/test_test_configurations.py b/test/tools/test_test_configurations.py new file mode 100644 index 000000000000..3dca70b7af53 --- /dev/null +++ b/test/tools/test_test_configurations.py @@ -0,0 +1,116 @@ +import os +import sys +import tempfile +import unittest +from types import SimpleNamespace +from unittest import mock + +ROOT = os.path.abspath(os.path.join(os.path.dirname(__file__), '..', '..')) +sys.path.insert(0, os.path.join(ROOT, 'tools')) +sys.path.insert(1, os.path.join(ROOT, 'test')) +import test as runner +import testpy + + +class TestConfigurationTest(unittest.TestCase): + def setUp(self): + directory = tempfile.TemporaryDirectory() + self.addCleanup(directory.cleanup) + self.root = directory.name + self.context = runner.Context( + ROOT, False, sys.executable, [], False, 5, lambda args: args, + False, False, 1, False) + for filename in ['test-example.js', 'example/test.js', + 'example/other.js', 'other/test.js']: + path = os.path.join(self.root, filename) + os.makedirs(os.path.dirname(path), exist_ok=True) + with open(path, 'w', encoding='utf8') as source: + source.write('') + + def module(self, suite): + return runner.get_module('testcfg', os.path.join(ROOT, 'test', suite)) + + def cases(self, config, suite): + return config.ListTests([suite], runner.SplitPath(suite), 'none', 'release') + + def test_base_case_and_simple_configuration_default_to_parallel(self): + case = runner.TestCase(self.context, ['example', 'test-example'], 'none', 'release') + self.assertTrue(case.parallel) + config = testpy.SimpleTestConfiguration(self.context, self.root, 'example') + cases = self.cases(config, 'example') + self.assertEqual(len(cases), 1) + self.assertTrue(cases[0].parallel) + + def test_serial_configuration_explicitly_disables_parallel_execution(self): + config = testpy.SerialTestConfiguration(self.context, self.root, 'example') + cases = self.cases(config, 'example') + self.assertEqual(len(cases), 1) + self.assertFalse(cases[0].parallel) + + def test_generic_and_serial_addon_configurations_preserve_discovery(self): + expected = [('example', 'other'), ('example', 'test'), ('other', 'test')] + for configuration, parallel in [(testpy.AddonTestConfiguration, True), + (testpy.SerialAddonTestConfiguration, False)]: + with self.subTest(configuration=configuration.__name__): + cases = self.cases(configuration(self.context, self.root, 'example'), 'example') + self.assertEqual(sorted(tuple(case.path[1:]) for case in cases), expected) + self.assertTrue(all(case.parallel == parallel for case in cases)) + + def test_named_suites_explicitly_remain_serial(self): + for suite in ['pummel', 'benchmark', 'known_issues', 'internet', + 'sequential', 'wasm-allocation', 'abort', + 'addons', 'js-native-api', 'node-api']: + with self.subTest(suite=suite): + config = self.module(suite).GetConfiguration(self.context, self.root) + cases = self.cases(config, suite) + self.assertTrue(cases) + self.assertTrue(all(not case.parallel for case in cases)) + + def test_other_suites_use_parallel_default(self): + for suite in ['parallel', 'async-hooks', 'client-proxy', 'doctool', + 'embedding', 'es-module', 'ffi', 'module-hooks', 'report', + 'sqlite', 'test426', 'test-runner', 'tick-processor', + 'trace_events', 'v8-updates', 'wasi']: + with self.subTest(suite=suite): + config = self.module(suite).GetConfiguration(self.context, self.root) + cases = self.cases(config, suite) + self.assertTrue(cases) + self.assertTrue(all(case.parallel for case in cases)) + + def test_known_issues_preserves_negative_context_without_changing_caller(self): + config = self.module('known_issues').GetConfiguration(self.context, self.root) + case = self.cases(config, 'known_issues')[0] + self.assertFalse(self.context.expect_fail) + self.assertTrue(case.IsNegative()) + self.assertIsNot(case.context, self.context) + self.assertFalse(case.parallel) + + def test_abort_retains_core_dump_suppression(self): + config = self.module('abort').GetConfiguration(self.context, self.root) + case = self.cases(config, 'abort')[0] + self.assertTrue(case.disable_core_files) + self.assertFalse(case.parallel) + + def test_skipped_wpt_wrapper_retains_serial_configuration(self): + config = self.module('wpt').GetConfiguration(self.context, self.root) + with mock.patch.object(config, '_Discover', return_value=None): + cases = self.cases(config, 'wpt') + self.assertEqual(len(cases), 1) + self.assertEqual(cases[0].path, ['wpt', 'test-example']) + self.assertFalse(cases[0].parallel) + + def test_sea_preserves_disk_space_gate_and_missing_binary_fallback(self): + sea = self.module('sea') + for exists, free, parallel in [(False, 1000, False), (True, 599, False), + (True, 600, True), (True, 1000, True)]: + with self.subTest(exists=exists, free=free): + with mock.patch.object(sea.os.path, 'isfile', return_value=exists), \ + mock.patch.object(sea.os.path, 'getsize', return_value=100), \ + mock.patch.object(sea.multiprocessing, 'cpu_count', return_value=2), \ + mock.patch.object(sea.shutil, 'disk_usage', return_value=SimpleNamespace(free=free)): + config = sea.GetConfiguration(self.context, self.root) + self.assertEqual(self.cases(config, 'sea')[0].parallel, parallel) + + +if __name__ == '__main__': + unittest.main() diff --git a/test/tools/test_test_runner.py b/test/tools/test_test_runner.py new file mode 100644 index 000000000000..29ffb93e8775 --- /dev/null +++ b/test/tools/test_test_runner.py @@ -0,0 +1,444 @@ +import contextlib +import io +import os +import sys +import threading +import unittest +from unittest import mock + +ROOT = os.path.abspath(os.path.join(os.path.dirname(__file__), '..', '..')) +sys.path.insert(0, os.path.join(ROOT, 'tools')) +import test as runner + + +WAIT_TIMEOUT = 10 + + +class SchedulerCase(runner.TestCase): + def __init__(self, name, action=None, suite='sequential', parallel=False, + arch='none', mode='release'): + super().__init__(None, [suite, name], arch, mode) + self.parallel = parallel + self.action = action + self.outcomes = {runner.PASS} + self.calls = 0 + + def IsNegative(self): + return False + + def Run(self): + self.calls += 1 + output = self.action(self) if self.action else None + if output is None: + output = runner.CommandOutput(0, False, '', '') + return runner.TestOutput(self, ['node', '/'.join(self.path)], output, False) + + +class SchedulerProgress(runner.ProgressIndicator): + def __init__(self, cases, flaky_tests_mode=runner.RUN, measure_flakiness=0, + reported=None): + super().__init__(cases, flaky_tests_mode, measure_flakiness) + self.started = [] + self.completed = [] + self.reported = reported or {} + + def Starting(self): + pass + + def Done(self): + pass + + def AboutToRun(self, case): + self.started.append(case) + + def HasRun(self, output): + self.completed.append(output.test) + event = self.reported.get(output.test) + if event: + event.set() + + +class SchedulerTest(unittest.TestCase): + def wait_for(self, event): + self.assertTrue(event.wait(WAIT_TIMEOUT), 'test runner did not make progress') + + @contextlib.contextmanager + def running(self, progress, tasks=2, release=()): + finished = threading.Event() + result = {} + errors = [] + + def run(): + try: + result.update(progress.Run(tasks)) + except BaseException as error: + errors.append(error) + finally: + finished.set() + + thread = threading.Thread(target=run, daemon=True) + thread.start() + + def complete(): + self.wait_for(finished) + if errors: + raise errors[0] + return result + + try: + yield complete + finally: + progress.Shutdown() + for event in release: + event.set() + thread.join(WAIT_TIMEOUT) + self.assertFalse(thread.is_alive(), 'test runner did not shut down') + + def test_distinct_subsystems_overlap_and_bypass_blocked_head(self): + first_started = threading.Event() + other_started = threading.Event() + release = threading.Event() + first_finished = threading.Event() + + def first(case): + first_started.set() + self.wait_for(other_started) + self.wait_for(release) + first_finished.set() + + def second(case): + self.assertTrue(first_finished.is_set()) + + def other(case): + self.wait_for(first_started) + other_started.set() + self.wait_for(release) + + cases = [SchedulerCase('test-net-first', first), + SchedulerCase('test-net-second', second), + SchedulerCase('test-http-first', other)] + progress = SchedulerProgress(cases) + with self.running(progress, release=[release]) as complete: + self.wait_for(other_started) + self.assertNotIn(cases[1], progress.started) + release.set() + self.assertTrue(complete()['allPassed']) + self.assertEqual(progress.remaining, 0) + self.assertEqual([case.calls for case in cases], [1, 1, 1]) + + def test_repeats_and_build_variants_preserve_subsystem_order(self): + entered = [] + active = set() + lock = threading.Lock() + first_started = threading.Event() + other_started = threading.Event() + release = threading.Event() + + def action(case): + with lock: + self.assertNotIn('net', active) + active.add('net') + entered.append(case) + if case is cases[0]: + first_started.set() + self.wait_for(other_started) + self.wait_for(release) + with lock: + active.remove('net') + + def other(case): + self.wait_for(first_started) + other_started.set() + + cases = [SchedulerCase('test-net-repeat', action), + SchedulerCase('test-net-repeat', action), + SchedulerCase('test-net-repeat', action, mode='debug'), + SchedulerCase('test-net-repeat', action, arch='arm64')] + progress = SchedulerProgress(cases + [SchedulerCase('test-http-other', other)]) + with self.running(progress, release=[release]) as complete: + self.wait_for(other_started) + self.assertEqual(entered, cases[:1]) + release.set() + self.assertTrue(complete()['allPassed']) + self.assertEqual(entered, cases) + + def test_single_job_uses_calling_thread_and_preserves_order(self): + calling_thread = threading.current_thread() + calls = [] + + def action(case): + self.assertIs(threading.current_thread(), calling_thread) + calls.append(case) + + cases = [SchedulerCase('test-net-first', action), + SchedulerCase('test-http-first', action), + SchedulerCase('test-net-second', action)] + progress = SchedulerProgress(cases) + with mock.patch.object(runner.threading, 'Thread') as worker: + self.assertTrue(progress.Run(1)['allPassed']) + worker.assert_not_called() + self.assertEqual(calls, cases) + self.assertEqual([case.thread_id for case in cases], [0, 0, 0]) + + def test_parallel_then_legacy_then_grouped_sequential_phases(self): + parallel_started = [threading.Event(), threading.Event()] + parallel_finished = [threading.Event(), threading.Event()] + release_parallel = threading.Event() + legacy_started = threading.Event() + legacy_finished = threading.Event() + release_legacy = threading.Event() + legacy_calls = [] + + def parallel(index): + def action(case): + parallel_started[index].set() + self.wait_for(release_parallel) + parallel_finished[index].set() + return action + + def legacy(case): + self.assertTrue(all(event.is_set() for event in parallel_finished)) + self.assertEqual(case.thread_id, 0) + legacy_calls.append(case) + if len(legacy_calls) == 1: + legacy_started.set() + self.wait_for(release_legacy) + else: + legacy_finished.set() + + def sequential(case): + self.assertTrue(legacy_finished.is_set()) + + legacy_cases = [SchedulerCase('test-addon-first', legacy, suite='addons'), + SchedulerCase('test-addon-second', legacy, suite='node-api')] + cases = [SchedulerCase('test-net-first', sequential), legacy_cases[0], + SchedulerCase('test-parallel-first', parallel(0), suite='parallel', parallel=True), + legacy_cases[1], + SchedulerCase('test-parallel-second', parallel(1), suite='parallel', parallel=True), + SchedulerCase('test-http-first', sequential)] + progress = SchedulerProgress(cases) + with self.running(progress, release=[release_parallel, release_legacy]) as complete: + for event in parallel_started: + self.wait_for(event) + release_parallel.set() + self.wait_for(legacy_started) + release_legacy.set() + self.assertTrue(complete()['allPassed']) + self.assertEqual(legacy_calls, legacy_cases) + + def test_assigns_unique_serial_ids_and_bounded_worker_ids(self): + cases = [SchedulerCase('test-net-%d' % index) for index in range(4)] + cases += [SchedulerCase('test-http-%d' % index) for index in range(4)] + cases += [SchedulerCase('test-parallel', suite='parallel', parallel=True), + SchedulerCase('test-legacy', suite='pummel')] + progress = SchedulerProgress(cases) + self.assertTrue(progress.Run(3)['allPassed']) + self.assertEqual(sorted(case.serial_id for case in cases), list(range(len(cases)))) + self.assertTrue(all(0 <= case.thread_id < 3 for case in cases)) + self.assertEqual(len(progress.completed), len(cases)) + self.assertTrue(all(case.duration is not None for case in cases)) + + def check_retries(self, measure_flakiness): + retry_started = threading.Event() + other_reported = threading.Event() + retry_finished = threading.Event() + + def retry(case): + if case.calls == 1: + return runner.CommandOutput(1, False, '', '') + retry_started.set() + self.wait_for(other_reported) + if case.calls == (3 if measure_flakiness else 2): + retry_finished.set() + + def follower(case): + self.assertTrue(retry_finished.is_set()) + + def other(case): + self.wait_for(retry_started) + + retried = SchedulerCase('test-net-retried', retry) + if not measure_flakiness: + retried.outcomes.add(runner.FLAKY) + unrelated = SchedulerCase('test-http-unrelated', other) + progress = SchedulerProgress( + [retried, SchedulerCase('test-net-follower', follower), unrelated], + flaky_tests_mode=runner.RUN if measure_flakiness else runner.KEEP_RETRYING, + measure_flakiness=measure_flakiness, reported={unrelated: other_reported}) + with contextlib.redirect_stdout(io.StringIO()): + with self.running(progress, release=[other_reported]) as complete: + result = complete() + self.assertEqual(retried.calls, 3 if measure_flakiness else 2) + self.assertEqual(result['allPassed'], not bool(measure_flakiness)) + self.assertEqual(progress.remaining, 0) + + def test_flaky_retry_retains_gate_and_allows_other_progress(self): + self.check_retries(0) + + def test_flakiness_measurement_retains_gate_and_allows_other_progress(self): + self.check_retries(2) + + def test_failed_crashed_and_timed_out_cases_release_subsystem(self): + for exit_code, timed_out, crashed in [(1, False, 0), (-9, False, 1), (-15, True, 0)]: + with self.subTest(exit_code=exit_code, timed_out=timed_out): + def fail(case): + return runner.CommandOutput(exit_code, timed_out, '', '') + + cases = [SchedulerCase('test-net-failed', fail), + SchedulerCase('test-net-following'), SchedulerCase('test-http-other')] + progress = SchedulerProgress(cases) + with mock.patch.object(runner.utils, 'IsWindows', return_value=False): + result = progress.Run(2) + self.assertFalse(result['allPassed']) + self.assertEqual(len(result['failed']), 1) + self.assertEqual(progress.crashed, crashed) + self.assertEqual(progress.remaining, 0) + self.assertEqual([case.calls for case in cases], [1, 1, 1]) + + def check_worker_error(self, error): + worker_started = threading.Event() + worker_failed = threading.Event() + waiter_started = threading.Event() + release = threading.Event() + + def action(case): + if case.thread_id: + worker_started.set() + self.wait_for(release) + worker_failed.set() + raise type(error)(str(error)) + self.wait_for(worker_failed) + + progress = SchedulerProgress([SchedulerCase('test-net-first', action), + SchedulerCase('test-net-following'), + SchedulerCase('test-http-first', action), + SchedulerCase('test-http-following')]) + original_wait = progress.condition.wait + + def wait(*args, **kwargs): + waiter_started.set() + return original_wait(*args, **kwargs) + + with mock.patch.object(progress.condition, 'wait', side_effect=wait): + with self.running(progress, tasks=3, release=[release, worker_failed]) as complete: + self.wait_for(worker_started) + self.wait_for(waiter_started) + release.set() + if isinstance(error, IOError): + self.assertFalse(complete()['allPassed']) + else: + with self.assertRaisesRegex(type(error), str(error)): + complete() + self.assertTrue(progress.shutdown_event.is_set()) + self.assertFalse(progress.running_subsystems) + + def test_shutdown_wakes_same_subsystem_waiters(self): + first_started = threading.Event() + waiter_started = threading.Event() + release = threading.Event() + + def first(case): + first_started.set() + self.wait_for(release) + + cases = [SchedulerCase('test-net-first', first), SchedulerCase('test-net-following')] + progress = SchedulerProgress(cases) + original_wait = progress.condition.wait + + def wait(*args, **kwargs): + waiter_started.set() + return original_wait(*args, **kwargs) + + with mock.patch.object(progress.condition, 'wait', side_effect=wait): + with self.running(progress, release=[release]) as complete: + self.wait_for(first_started) + self.wait_for(waiter_started) + progress.Shutdown() + release.set() + self.assertFalse(complete()['allPassed']) + self.assertEqual([case.calls for case in cases], [1, 0]) + self.assertFalse(progress.running_subsystems) + + def test_interrupt_during_join_waits_for_worker_cleanup(self): + for interrupt in [KeyboardInterrupt, SystemExit]: + with self.subTest(interrupt=interrupt.__name__): + worker_started = threading.Event() + interrupted = threading.Event() + joined_again = threading.Event() + release = threading.Event() + done = threading.Event() + workers = [] + + def action(case): + workers.append(threading.current_thread()) + worker_started.set() + self.wait_for(release) + + progress = SchedulerProgress([SchedulerCase('test-net-first', action)]) + original_run_single = progress.RunSingle + original_join = threading.Thread.join + original_is_alive = threading.Thread.is_alive + + def run_single(queue, thread_id): + if thread_id == 0 and queue is progress.sequential_queue: + self.wait_for(worker_started) + original_run_single(queue, thread_id) + + def join(thread, timeout=None): + if thread in workers: + if not interrupted.is_set(): + interrupted.set() + raise interrupt() + joined_again.set() + original_join(thread, timeout) + + def is_alive(thread): + # CPython 3.12 can report stopped after SIGINT interrupts join. + if thread in workers and interrupted.is_set(): + return False + return original_is_alive(thread) + + with mock.patch.object(progress, 'RunSingle', side_effect=run_single), \ + mock.patch.object(progress, 'Done', side_effect=done.set), \ + mock.patch.object(threading.Thread, 'join', join), \ + mock.patch.object(threading.Thread, 'is_alive', is_alive): + with self.running(progress, release=[release]) as complete: + self.wait_for(interrupted) + self.wait_for(joined_again) + self.assertFalse(done.is_set()) + self.assertEqual(progress.running_subsystems, {'net'}) + release.set() + self.assertFalse(complete()['allPassed']) + self.assertTrue(done.is_set()) + self.assertFalse(progress.running_subsystems) + + def test_worker_start_failure_does_not_wait_for_unstarted_thread(self): + case = SchedulerCase('test-net-first') + progress = SchedulerProgress([case]) + with mock.patch.object(threading.Thread, 'start', side_effect=RuntimeError('cannot start worker')): + with self.assertRaisesRegex(RuntimeError, 'cannot start worker'): + progress.Run(2) + self.assertEqual(case.calls, 0) + self.assertTrue(progress.shutdown_event.is_set()) + self.assertFalse(progress.running_subsystems) + + def test_worker_ioerror_aborts_run(self): + self.check_worker_error(IOError('scheduler fixture interrupted')) + + def test_worker_exception_is_propagated_to_caller(self): + self.check_worker_error(RuntimeError('scheduler fixture failed')) + + def test_progress_callback_exceptions_release_subsystem(self): + for callback in ['AboutToRun', 'HasRun']: + with self.subTest(callback=callback): + progress = SchedulerProgress([SchedulerCase('test-net-first'), + SchedulerCase('test-net-following')]) + with mock.patch.object(progress, callback, side_effect=RuntimeError('reporter failed')): + with self.running(progress) as complete: + with self.assertRaisesRegex(RuntimeError, 'reporter failed'): + complete() + self.assertFalse(progress.running_subsystems) + + +if __name__ == '__main__': + unittest.main() diff --git a/test/trace_events/testcfg.py b/test/trace_events/testcfg.py index e62bbff517f7..c2ed27e62a3d 100644 --- a/test/trace_events/testcfg.py +++ b/test/trace_events/testcfg.py @@ -3,4 +3,4 @@ import testpy def GetConfiguration(context, root): - return testpy.ParallelTestConfiguration(context, root, 'trace_events') + return testpy.SimpleTestConfiguration(context, root, 'trace_events') diff --git a/test/v8-updates/testcfg.py b/test/v8-updates/testcfg.py index cec2589f9b4c..a448702a0b51 100644 --- a/test/v8-updates/testcfg.py +++ b/test/v8-updates/testcfg.py @@ -3,4 +3,4 @@ import testpy def GetConfiguration(context, root): - return testpy.ParallelTestConfiguration(context, root, 'v8-updates') + return testpy.SimpleTestConfiguration(context, root, 'v8-updates') diff --git a/test/wasi/testcfg.py b/test/wasi/testcfg.py index ec6cbc5fe3dc..4e0a7ce2156f 100644 --- a/test/wasi/testcfg.py +++ b/test/wasi/testcfg.py @@ -3,4 +3,4 @@ import testpy def GetConfiguration(context, root): - return testpy.ParallelTestConfiguration(context, root, 'wasi') + return testpy.SimpleTestConfiguration(context, root, 'wasi') diff --git a/test/wasm-allocation/testcfg.py b/test/wasm-allocation/testcfg.py index 4962550b4b69..d050b645acca 100644 --- a/test/wasm-allocation/testcfg.py +++ b/test/wasm-allocation/testcfg.py @@ -3,4 +3,5 @@ import testpy def GetConfiguration(context, root): - return testpy.SimpleTestConfiguration(context, root, 'wasm-allocation') + # TODO: Replace preexec_fn memory limits before parallelizing this suite. + return testpy.SerialTestConfiguration(context, root, 'wasm-allocation') diff --git a/test/wpt/testcfg.py b/test/wpt/testcfg.py index 2f3097da96df..80ed22f7d3f3 100644 --- a/test/wpt/testcfg.py +++ b/test/wpt/testcfg.py @@ -17,6 +17,7 @@ def __init__(self, path, file, arch, mode, context, config, group, serial): super(WPTTestCase, self).__init__( path, file, arch, mode, context, config, config.additional_flags) self.group = group + # TODO: Audit serial groups for process isolation and timer sensitivity. self.parallel = not serial def GetName(self): @@ -44,7 +45,7 @@ def GetRunConfiguration(self): return configuration -class WPTTestConfiguration(testpy.SimpleTestConfiguration): +class WPTTestConfiguration(testpy.SerialTestConfiguration): def __init__(self, context, root): super(WPTTestConfiguration, self).__init__(context, root, 'wpt') self.manifests = {} diff --git a/tools/test.py b/tools/test.py index 90df2cfb2fac..f3deee95744e 100755 --- a/tools/test.py +++ b/tools/test.py @@ -103,12 +103,17 @@ def __init__(self, cases, flaky_tests_mode, measure_flakiness): self.flaky_tests_mode = flaky_tests_mode self.measure_flakiness = measure_flakiness self.parallel_queue = Queue(len(cases)) - self.sequential_queue = Queue(len(cases)) + self.serial_queue = Queue(len(cases)) + self.sequential_queue = [] + self.running_subsystems = set() + self.condition = threading.Condition() for case in cases: if case.parallel: self.parallel_queue.put_nowait(case) + elif case.path[0] == 'sequential': + self.sequential_queue.append(case) else: - self.sequential_queue.put_nowait(case) + self.serial_queue.put_nowait(case) self.succeeded = 0 self.remaining = len(cases) self.total = len(cases) @@ -117,6 +122,7 @@ def __init__(self, cases, flaky_tests_mode, measure_flakiness): self.crashed = 0 self.lock = threading.Lock() self.shutdown_event = threading.Event() + self.worker_errors = [] def GetFailureOutput(self, failure): output = [] @@ -150,77 +156,165 @@ def PrintFailureHeader(self, test): def Run(self, tasks) -> Dict: self.Starting() - threads = [] - # Spawn N-1 threads and then use this thread as the last one. - # That way -j1 avoids threading altogether which is a nice fallback - # in case of threading problems. - for i in range(tasks - 1): - thread = threading.Thread(target=self.RunSingle, args=[True, i + 1]) - threads.append(thread) - thread.start() try: - self.RunSingle(False, 0) - # Wait for the remaining threads - for thread in threads: - # Use a timeout so that signals (ctrl-c) will be processed. - thread.join(timeout=1000000) + self.RunPhase(self.parallel_queue, tasks) + self.RunSingle(self.serial_queue, 0) + self.RunPhase(self.sequential_queue, tasks) except (KeyboardInterrupt, SystemExit): - self.shutdown_event.set() + self.Shutdown() except Exception: - # If there's an exception we schedule an interruption for any - # remaining threads. - self.shutdown_event.set() - # ...and then reraise the exception to bail out + self.Shutdown() raise self.Done() return { - 'allPassed': not self.failed and not self.shutdown_event.is_set(), + 'allPassed': not self.failed and not self.shutdown_event.is_set() and self.remaining == 0, 'failed': self.failed, } - def RunSingle(self, parallel, thread_id): + def Shutdown(self): + with self.condition: + self.shutdown_event.set() + self.condition.notify_all() + + def RunPhase(self, queue, tasks): + if self.shutdown_event.is_set(): + return + if queue is self.sequential_queue: + empty = not queue + else: + empty = queue.empty() + if empty: + return + threads = [] + # Spawn N-1 threads and then use this thread as the last one. + # That way -j1 avoids threading altogether which is a nice fallback + # in case of threading problems. + try: + for i in range(tasks - 1): + finished = threading.Event() + thread = threading.Thread(target=self.RunWorker, args=[queue, i + 1, finished]) + threads.append((thread, finished)) + thread.start() + self.RunSingle(queue, 0) + except BaseException: + self.Shutdown() + raise + finally: + for thread, finished in threads: + if thread.ident is None: + continue + # Use a timeout so that signals (ctrl-c) will be processed. + # Interrupted joins can mark a live thread stopped on some Python versions. + while not finished.is_set(): + try: + thread.join(timeout=0.1) + finished.wait(timeout=0.1) + except (KeyboardInterrupt, SystemExit): # noqa: PERF203 + self.Shutdown() + if self.worker_errors: + _, error, traceback = self.worker_errors[0] + raise error.with_traceback(traceback) + + def RunWorker(self, queue, thread_id, finished): + try: + self.RunSingle(queue, thread_id) + except BaseException: + with self.lock: + self.worker_errors.append(sys.exc_info()) + self.Shutdown() + finally: + finished.set() + + def GetSequentialTest(self): + with self.condition: + while not self.shutdown_event.is_set(): + for index, case in enumerate(self.sequential_queue): + # The first component after test- names the subsystem. Skip busy + # subsystems without changing the order of their pending tests. + subsystem = case.path[-1].split('-', 2)[1] + if subsystem not in self.running_subsystems: + self.running_subsystems.add(subsystem) + return self.sequential_queue.pop(index) + if not self.sequential_queue: + return None + self.condition.wait() + return None + + def RunSingle(self, queue, thread_id): while not self.shutdown_event.is_set(): - try: - test = self.parallel_queue.get_nowait() - except Empty: - if parallel: + sequential = queue is self.sequential_queue + if sequential: + case = self.GetSequentialTest() + if case is None: return + else: try: - test = self.sequential_queue.get_nowait() + case = queue.get_nowait() except Empty: return - case = test + try: + if self.shutdown_event.is_set(): + return + self.RunCase(case, thread_id) + except IOError: + self.Shutdown() + return + except BaseException: + self.Shutdown() + raise + finally: + if sequential: + with self.condition: + self.running_subsystems.remove(case.path[-1].split('-', 2)[1]) + self.condition.notify_all() + + def RunCase(self, case, thread_id): + with self.lock: case.thread_id = thread_id - self.lock.acquire() case.serial_id = self.serial_id self.serial_id += 1 self.AboutToRun(case) - self.lock.release() - try: - start = datetime.now() + start = datetime.now() + output = case.Run() + # SmartOS has a bug that causes unexpected ECONNREFUSED errors. + # See https://smartos.org/bugview/OS-2767 + # If ECONNREFUSED on SmartOS, retry the test one time. + if (output.UnexpectedOutput() and + sys.platform == 'sunos5' and + 'ECONNREFUSED' in output.output.stderr): output = case.Run() - # SmartOS has a bug that causes unexpected ECONNREFUSED errors. - # See https://smartos.org/bugview/OS-2767 - # If ECONNREFUSED on SmartOS, retry the test one time. - if (output.UnexpectedOutput() and - sys.platform == 'sunos5' and - 'ECONNREFUSED' in output.output.stderr): - output = case.Run() - output.diagnostic.append('ECONNREFUSED received, test retried') - case.duration = (datetime.now() - start) - except IOError: - return - if self.shutdown_event.is_set(): - return - self.lock.acquire() - if output.UnexpectedOutput(): - if FLAKY in output.test.outcomes and self.flaky_tests_mode == DONTCARE: + output.diagnostic.append('ECONNREFUSED received, test retried') + case.duration = (datetime.now() - start) + if self.shutdown_event.is_set(): + return + + unexpected = output.UnexpectedOutput() + flaky = FLAKY in output.test.outcomes + retry_passed = False + measured_failures = None + if unexpected and flaky and self.flaky_tests_mode == KEEP_RETRYING: + for _ in range(99): + if self.shutdown_event.is_set(): + return + if not case.Run().UnexpectedOutput(): + retry_passed = True + break + elif unexpected and not (flaky and self.flaky_tests_mode == DONTCARE) and self.measure_flakiness: + measured_failures = 1 + for _ in range(self.measure_flakiness): + if self.shutdown_event.is_set(): + return + measured_failures += bool(case.Run().UnexpectedOutput()) + + if self.shutdown_event.is_set(): + return + with self.lock: + if unexpected: + if flaky and self.flaky_tests_mode == DONTCARE: self.flaky_failed.append(output) - elif FLAKY in output.test.outcomes and self.flaky_tests_mode == KEEP_RETRYING: - for _ in range(99): - if not case.Run().UnexpectedOutput(): - self.flaky_failed.append(output) - break + elif flaky and self.flaky_tests_mode == KEEP_RETRYING: + if retry_passed: + self.flaky_failed.append(output) else: # If after 100 tries, the test is not passing, it's not flaky. self.failed.append(output) @@ -228,15 +322,13 @@ def RunSingle(self, parallel, thread_id): self.failed.append(output) if output.HasCrashed(): self.crashed += 1 - if self.measure_flakiness: - outputs = [case.Run() for _ in range(self.measure_flakiness)] + if measured_failures is not None: # +1s are there because the test already failed once at this point. - print(" failed %d out of %d" % (len([i for i in outputs if i.UnexpectedOutput()]) + 1, self.measure_flakiness + 1)) + print(" failed %d out of %d" % (measured_failures, self.measure_flakiness + 1)) else: self.succeeded += 1 self.remaining -= 1 self.HasRun(output) - self.lock.release() def EscapeCommand(command): @@ -563,7 +655,7 @@ def __init__(self, context, path, arch, mode): self.duration = None self.arch = arch self.mode = mode - self.parallel = False + self.parallel = True self.disable_core_files = False self.max_virtual_memory = None self.serial_id = 0 From e3706bb7f2f8ea13d13ee2b72b04b4789de526c4 Mon Sep 17 00:00:00 2001 From: Filip Skokan Date: Thu, 1 Oct 2026 14:13:52 +0200 Subject: [PATCH 02/10] test: size SEA disk guard for requested workers Use the resolved test runner worker count instead of the host CPU count when checking space for concurrent executable copies. Retain the serial fallback when the executable is missing or free space is insufficient. Signed-off-by: Filip Skokan Assisted-by: Codex --- test/sea/testcfg.py | 6 ++---- test/tools/test_test_configurations.py | 13 ++++++++----- tools/test.py | 2 ++ 3 files changed, 12 insertions(+), 9 deletions(-) diff --git a/test/sea/testcfg.py b/test/sea/testcfg.py index bc08b157add5..6cd439f18a40 100644 --- a/test/sea/testcfg.py +++ b/test/sea/testcfg.py @@ -1,4 +1,4 @@ -import sys, os, multiprocessing, shutil +import sys, os, shutil sys.path.append(os.path.join(os.path.dirname(__file__), '..')) import testpy @@ -18,12 +18,10 @@ def GetConfiguration(context, root): return testpy.SerialTestConfiguration(context, root, 'sea') # Get the size of the executable to decide whether we can run tests in parallel. - # TODO: Use the requested worker count when evaluating this disk-space guard. executable_size = os.path.getsize(vm) - num_cpus = multiprocessing.cpu_count() remaining_disk_space = shutil.disk_usage('.').free # Give it a bit of leeway by multiplying by 3. - if (executable_size * num_cpus * 3 > remaining_disk_space): + if (executable_size * context.jobs * 3 > remaining_disk_space): return testpy.SerialTestConfiguration(context, root, 'sea') return testpy.SimpleTestConfiguration(context, root, 'sea') diff --git a/test/tools/test_test_configurations.py b/test/tools/test_test_configurations.py index 3dca70b7af53..d8acbda8838c 100644 --- a/test/tools/test_test_configurations.py +++ b/test/tools/test_test_configurations.py @@ -99,14 +99,17 @@ def test_skipped_wpt_wrapper_retains_serial_configuration(self): self.assertEqual(cases[0].path, ['wpt', 'test-example']) self.assertFalse(cases[0].parallel) - def test_sea_preserves_disk_space_gate_and_missing_binary_fallback(self): + def test_sea_uses_requested_jobs_for_disk_space_gate(self): sea = self.module('sea') - for exists, free, parallel in [(False, 1000, False), (True, 599, False), - (True, 600, True), (True, 1000, True)]: - with self.subTest(exists=exists, free=free): + for exists, jobs, free, parallel in [(False, 1, 1000, False), + (True, 1, 299, False), + (True, 1, 300, True), + (True, 4, 1199, False), + (True, 4, 1200, True)]: + with self.subTest(exists=exists, jobs=jobs, free=free): + self.context.jobs = jobs with mock.patch.object(sea.os.path, 'isfile', return_value=exists), \ mock.patch.object(sea.os.path, 'getsize', return_value=100), \ - mock.patch.object(sea.multiprocessing, 'cpu_count', return_value=2), \ mock.patch.object(sea.shutil, 'disk_usage', return_value=SimpleNamespace(free=free)): config = sea.GetConfiguration(self.context, self.root) self.assertEqual(self.cases(config, 'sea')[0].parallel, parallel) diff --git a/tools/test.py b/tools/test.py index f3deee95744e..b77439c59c94 100755 --- a/tools/test.py +++ b/tools/test.py @@ -1063,6 +1063,7 @@ def __init__(self, workspace, verbose, vm, args, expect_fail, self.store_unexpected_output = store_unexpected_output self.repeat = repeat self.abort_on_timeout = abort_on_timeout + self.jobs = 1 self.v8_enable_inspector = True self.node_has_crypto = True self.node_has_ffi = True @@ -1797,6 +1798,7 @@ def Main(): options.store_unexpected_output, options.repeat, options.abort_on_timeout) + context.jobs = options.j # Remember the primary mode requested on the CLI so suites can reuse it when # they need to probe for a binary outside of the normal test runner flow. for requested_mode in options.mode: From 6b0c6d99525e26d985c9678503c453f06b31860d Mon Sep 17 00:00:00 2001 From: Filip Skokan Date: Thu, 1 Oct 2026 14:18:17 +0200 Subject: [PATCH 03/10] test: run abort tests with parallel-safe limits Set core limits in a launcher process that execs the test command instead of using preexec_fn in the threaded runner. Preserve process identity, timeouts and exit status, then allow the abort suite to run in parallel. Signed-off-by: Filip Skokan Assisted-by: Codex --- test/README.md | 4 +- test/abort/testcfg.py | 1 - test/testpy/__init__.py | 2 +- test/tools/test_test_configurations.py | 6 +-- test/tools/test_test_resource_limits.py | 70 +++++++++++++++++++++++++ tools/test-resource-limits.py | 19 +++++++ tools/test.py | 22 +++----- 7 files changed, 102 insertions(+), 22 deletions(-) create mode 100644 test/tools/test_test_resource_limits.py create mode 100644 tools/test-resource-limits.py diff --git a/test/README.md b/test/README.md index 37d54c544c98..4f148b15c84d 100644 --- a/test/README.md +++ b/test/README.md @@ -41,8 +41,8 @@ For the tests to run on Windows, be sure to clone Node.js source code with the Tests run in parallel by default. Suite configurations opt out with `SerialTestConfiguration`, or `SerialAddonTestConfiguration` for addon layouts. The `addons`, `benchmark`, `internet`, `js-native-api`, `known_issues`, `node-api`, -and `pummel` suites explicitly run serially. The `abort` and `wasm-allocation` -suites also remain serial; WPT's group settings and SEA's disk-space guard retain +and `pummel` suites explicitly run serially. The `wasm-allocation` +suite also remains serial; WPT's group settings and SEA's disk-space guard retain their existing scheduling. The test runner finishes parallel tests first, followed by serial suites, then diff --git a/test/abort/testcfg.py b/test/abort/testcfg.py index 40b1ab429175..e509d0453c40 100644 --- a/test/abort/testcfg.py +++ b/test/abort/testcfg.py @@ -3,5 +3,4 @@ import testpy def GetConfiguration(context, root): - # TODO: Replace preexec_fn core suppression before parallelizing this suite. return testpy.AbortTestConfiguration(context, root, 'abort') diff --git a/test/testpy/__init__.py b/test/testpy/__init__.py index 816b6437490b..8f56c3a28a3b 100644 --- a/test/testpy/__init__.py +++ b/test/testpy/__init__.py @@ -186,7 +186,7 @@ class SerialAddonTestConfiguration(AddonTestConfiguration): parallel = False -class AbortTestConfiguration(SerialTestConfiguration): +class AbortTestConfiguration(SimpleTestConfiguration): def __init__(self, context, root, section, additional=None): super(AbortTestConfiguration, self).__init__(context, root, section, additional) diff --git a/test/tools/test_test_configurations.py b/test/tools/test_test_configurations.py index d8acbda8838c..faf48af8395d 100644 --- a/test/tools/test_test_configurations.py +++ b/test/tools/test_test_configurations.py @@ -58,7 +58,7 @@ def test_generic_and_serial_addon_configurations_preserve_discovery(self): def test_named_suites_explicitly_remain_serial(self): for suite in ['pummel', 'benchmark', 'known_issues', 'internet', - 'sequential', 'wasm-allocation', 'abort', + 'sequential', 'wasm-allocation', 'addons', 'js-native-api', 'node-api']: with self.subTest(suite=suite): config = self.module(suite).GetConfiguration(self.context, self.root) @@ -67,7 +67,7 @@ def test_named_suites_explicitly_remain_serial(self): self.assertTrue(all(not case.parallel for case in cases)) def test_other_suites_use_parallel_default(self): - for suite in ['parallel', 'async-hooks', 'client-proxy', 'doctool', + for suite in ['parallel', 'abort', 'async-hooks', 'client-proxy', 'doctool', 'embedding', 'es-module', 'ffi', 'module-hooks', 'report', 'sqlite', 'test426', 'test-runner', 'tick-processor', 'trace_events', 'v8-updates', 'wasi']: @@ -89,7 +89,7 @@ def test_abort_retains_core_dump_suppression(self): config = self.module('abort').GetConfiguration(self.context, self.root) case = self.cases(config, 'abort')[0] self.assertTrue(case.disable_core_files) - self.assertFalse(case.parallel) + self.assertTrue(case.parallel) def test_skipped_wpt_wrapper_retains_serial_configuration(self): config = self.module('wpt').GetConfiguration(self.context, self.root) diff --git a/test/tools/test_test_resource_limits.py b/test/tools/test_test_resource_limits.py new file mode 100644 index 000000000000..322ccff3de50 --- /dev/null +++ b/test/tools/test_test_resource_limits.py @@ -0,0 +1,70 @@ +import concurrent.futures +import json +import os +import sys +import unittest +from types import SimpleNamespace +from unittest import mock + +ROOT = os.path.abspath(os.path.join(os.path.dirname(__file__), '..', '..')) +sys.path.insert(0, os.path.join(ROOT, 'tools')) +import test as runner + +if sys.platform != 'win32': + import resource + + +@unittest.skipIf(sys.platform == 'win32', 'POSIX resource limits') +class ResourceLimitsTest(unittest.TestCase): + def setUp(self): + self.context = SimpleNamespace(verbose=False, suppress_dialogs=False, + abort_on_timeout=False) + + def test_core_limits_are_isolated_between_concurrent_processes(self): + parent_limits = resource.getrlimit(resource.RLIMIT_CORE) + code = '''import json, os, resource, sys +print(json.dumps([resource.getrlimit(resource.RLIMIT_CORE), + sys.argv[1:], os.environ['RESOURCE_LIMIT_TEST']])) +''' + + def run(index): + result = runner.Execute([sys.executable, '-c', code, 'a b', '--literal'], + self.context, timeout=5, + env={'RESOURCE_LIMIT_TEST': str(index)}, + disable_core_files=True) + self.assertEqual(result.exit_code, 0, result.stderr) + self.assertEqual(json.loads(result.stdout), [[0, 0], ['a b', '--literal'], str(index)]) + + with concurrent.futures.ThreadPoolExecutor(max_workers=4) as pool: + list(pool.map(run, range(8))) + self.assertEqual(resource.getrlimit(resource.RLIMIT_CORE), parent_limits) + + def test_wrapper_exec_preserves_the_test_process_id(self): + processes = [] + original_popen = runner.subprocess.Popen + + def popen(*args, **kwargs): + self.assertIsNone(kwargs.get('preexec_fn')) + process = original_popen(*args, **kwargs) + processes.append(process) + return process + + with mock.patch.object(runner.subprocess, 'Popen', side_effect=popen): + result = runner.Execute([sys.executable, '-c', 'import os; print(os.getpid())'], + self.context, timeout=5, disable_core_files=True) + self.assertEqual(result.exit_code, 0, result.stderr) + self.assertEqual(int(result.stdout), processes[0].pid) + + def test_wrapper_preserves_failures_and_timeouts(self): + failure = runner.Execute([sys.executable, '-c', 'raise SystemExit(7)'], + self.context, timeout=5, disable_core_files=True) + self.assertEqual(failure.exit_code, 7) + self.assertFalse(failure.timed_out) + timeout = runner.Execute([sys.executable, '-c', 'import time; time.sleep(60)'], + self.context, timeout=0.1, disable_core_files=True) + self.assertTrue(timeout.timed_out) + self.assertNotEqual(timeout.exit_code, 0) + + +if __name__ == '__main__': + unittest.main() diff --git a/tools/test-resource-limits.py b/tools/test-resource-limits.py new file mode 100644 index 000000000000..dddccb71cadf --- /dev/null +++ b/tools/test-resource-limits.py @@ -0,0 +1,19 @@ +import argparse +import os +import resource + +parser = argparse.ArgumentParser() +parser.add_argument('--disable-core-files', action='store_true') +parser.add_argument('command', nargs=argparse.REMAINDER) +options = parser.parse_args() +command = options.command +if command and command[0] == '--': + command = command[1:] +if not command: + parser.error('a test command is required') + +if options.disable_core_files: + resource.setrlimit(resource.RLIMIT_CORE, (0, 0)) + +# Replace this process so test signals, timeouts and exit codes reach the runner. +os.execvpe(command[0], command, os.environ) diff --git a/tools/test.py b/tools/test.py index b77439c59c94..ae2c82d12957 100755 --- a/tools/test.py +++ b/tools/test.py @@ -891,12 +891,9 @@ def Execute(args, context, timeout=None, env=None, disable_core_files=False, preexec_fn = None - def disableCoreFiles(): - import resource - resource.setrlimit(resource.RLIMIT_CORE, (0,0)) - if disable_core_files and not utils.IsWindows(): - preexec_fn = disableCoreFiles + args = [sys.executable, join(dirname(__file__), 'test-resource-limits.py'), + '--disable-core-files', '--'] + args if max_virtual_memory is not None and utils.GuessOS() == 'linux': def setMaxVirtualMemory(): @@ -904,14 +901,7 @@ def setMaxVirtualMemory(): resource.setrlimit(resource.RLIMIT_CORE, (0,0)) resource.setrlimit(resource.RLIMIT_AS, (max_virtual_memory,max_virtual_memory + 1)) - if preexec_fn is not None: - prev_preexec_fn = preexec_fn - def setResourceLimits(): - setMaxVirtualMemory() - prev_preexec_fn() - preexec_fn = setResourceLimits - else: - preexec_fn = setMaxVirtualMemory + preexec_fn = setMaxVirtualMemory (_process, exit_code, timed_out) = RunProcess( context, @@ -925,8 +915,10 @@ def setResourceLimits(): ) os.close(fd_out) os.close(fd_err) - output = open(outname, encoding='utf8').read() - errors = open(errname, encoding='utf8').read() + with open(outname, encoding='utf8') as output_file: + output = output_file.read() + with open(errname, encoding='utf8') as error_file: + errors = error_file.read() CheckedUnlink(outname) CheckedUnlink(errname) From a2db6c100d5360d02b7521cb574c7c21f86396f9 Mon Sep 17 00:00:00 2001 From: Filip Skokan Date: Thu, 1 Oct 2026 14:22:10 +0200 Subject: [PATCH 04/10] test: launch wasm tests with isolated memory limits Apply Linux RLIMIT_AS and core limits in the test launcher, preserving the existing soft and hard bounds without preexec_fn. Allow wasm-allocation tests to run concurrently while retaining their platform-specific skips. Signed-off-by: Filip Skokan Assisted-by: Codex --- test/README.md | 5 ++- test/tools/test_test_configurations.py | 4 +-- test/tools/test_test_resource_limits.py | 46 +++++++++++++++++++++++++ test/wasm-allocation/testcfg.py | 3 +- tools/test-resource-limits.py | 6 +++- tools/test.py | 17 ++++----- 6 files changed, 63 insertions(+), 18 deletions(-) diff --git a/test/README.md b/test/README.md index 4f148b15c84d..9b2f2ca7e93d 100644 --- a/test/README.md +++ b/test/README.md @@ -41,9 +41,8 @@ For the tests to run on Windows, be sure to clone Node.js source code with the Tests run in parallel by default. Suite configurations opt out with `SerialTestConfiguration`, or `SerialAddonTestConfiguration` for addon layouts. The `addons`, `benchmark`, `internet`, `js-native-api`, `known_issues`, `node-api`, -and `pummel` suites explicitly run serially. The `wasm-allocation` -suite also remains serial; WPT's group settings and SEA's disk-space guard retain -their existing scheduling. +and `pummel` suites explicitly run serially. WPT's group settings and SEA's +disk-space guard retain their existing scheduling. The test runner finishes parallel tests first, followed by serial suites, then `sequential` tests. In `sequential`, different subsystems can run concurrently diff --git a/test/tools/test_test_configurations.py b/test/tools/test_test_configurations.py index faf48af8395d..5d92883d7e41 100644 --- a/test/tools/test_test_configurations.py +++ b/test/tools/test_test_configurations.py @@ -58,7 +58,7 @@ def test_generic_and_serial_addon_configurations_preserve_discovery(self): def test_named_suites_explicitly_remain_serial(self): for suite in ['pummel', 'benchmark', 'known_issues', 'internet', - 'sequential', 'wasm-allocation', + 'sequential', 'addons', 'js-native-api', 'node-api']: with self.subTest(suite=suite): config = self.module(suite).GetConfiguration(self.context, self.root) @@ -70,7 +70,7 @@ def test_other_suites_use_parallel_default(self): for suite in ['parallel', 'abort', 'async-hooks', 'client-proxy', 'doctool', 'embedding', 'es-module', 'ffi', 'module-hooks', 'report', 'sqlite', 'test426', 'test-runner', 'tick-processor', - 'trace_events', 'v8-updates', 'wasi']: + 'trace_events', 'v8-updates', 'wasi', 'wasm-allocation']: with self.subTest(suite=suite): config = self.module(suite).GetConfiguration(self.context, self.root) cases = self.cases(config, suite) diff --git a/test/tools/test_test_resource_limits.py b/test/tools/test_test_resource_limits.py index 322ccff3de50..b8f3ed485075 100644 --- a/test/tools/test_test_resource_limits.py +++ b/test/tools/test_test_resource_limits.py @@ -1,6 +1,7 @@ import concurrent.futures import json import os +import runpy import sys import unittest from types import SimpleNamespace @@ -65,6 +66,51 @@ def test_wrapper_preserves_failures_and_timeouts(self): self.assertTrue(timeout.timed_out) self.assertNotEqual(timeout.exit_code, 0) + def test_memory_launcher_sets_both_limits_before_exec(self): + limit = 512 * 1024 * 1024 + helper = os.path.join(ROOT, 'tools', 'test-resource-limits.py') + command = [sys.executable, '-c', 'pass'] + argv = [helper, '--max-virtual-memory', str(limit), '--'] + command + with mock.patch.object(sys, 'argv', argv), \ + mock.patch.object(resource, 'setrlimit') as setrlimit, \ + mock.patch.object(os, 'execvpe') as execvpe: + runpy.run_path(helper, run_name='__main__') + self.assertEqual(setrlimit.call_args_list, + [mock.call(resource.RLIMIT_CORE, (0, 0)), + mock.call(resource.RLIMIT_AS, (limit, limit + 1))]) + execvpe.assert_called_once_with(command[0], command, os.environ) + + def test_memory_limits_use_launcher_only_on_linux(self): + command = [sys.executable, '-c', 'pass'] + for platform in ['linux', 'macos']: + with self.subTest(platform=platform), \ + mock.patch.object(runner.utils, 'GuessOS', return_value=platform), \ + mock.patch.object(runner, 'RunProcess', return_value=(None, 0, False)) as run: + runner.Execute(command, self.context, max_virtual_memory=123456) + actual = run.call_args.kwargs + self.assertNotIn('preexec_fn', actual) + if platform == 'linux': + self.assertEqual(actual['args'], + [sys.executable, os.path.join(ROOT, 'tools', 'test-resource-limits.py'), + '--max-virtual-memory', '123456', '--'] + command) + else: + self.assertEqual(actual['args'], command) + + @unittest.skipUnless(sys.platform.startswith('linux'), 'Linux virtual memory limits') + def test_concurrent_memory_limits_are_inherited_by_tests(self): + parent_limits = resource.getrlimit(resource.RLIMIT_AS) + code = 'import json, resource; print(json.dumps(resource.getrlimit(resource.RLIMIT_AS)))' + + def run(limit): + result = runner.Execute([sys.executable, '-c', code], self.context, + timeout=5, max_virtual_memory=limit) + self.assertEqual(result.exit_code, 0, result.stderr) + self.assertEqual(json.loads(result.stdout), [limit, limit + 1]) + + with concurrent.futures.ThreadPoolExecutor(max_workers=4) as pool: + list(pool.map(run, [512 * 1024 * 1024, 768 * 1024 * 1024] * 4)) + self.assertEqual(resource.getrlimit(resource.RLIMIT_AS), parent_limits) + if __name__ == '__main__': unittest.main() diff --git a/test/wasm-allocation/testcfg.py b/test/wasm-allocation/testcfg.py index d050b645acca..4962550b4b69 100644 --- a/test/wasm-allocation/testcfg.py +++ b/test/wasm-allocation/testcfg.py @@ -3,5 +3,4 @@ import testpy def GetConfiguration(context, root): - # TODO: Replace preexec_fn memory limits before parallelizing this suite. - return testpy.SerialTestConfiguration(context, root, 'wasm-allocation') + return testpy.SimpleTestConfiguration(context, root, 'wasm-allocation') diff --git a/tools/test-resource-limits.py b/tools/test-resource-limits.py index dddccb71cadf..aee96b884c38 100644 --- a/tools/test-resource-limits.py +++ b/tools/test-resource-limits.py @@ -4,6 +4,7 @@ parser = argparse.ArgumentParser() parser.add_argument('--disable-core-files', action='store_true') +parser.add_argument('--max-virtual-memory', type=int) parser.add_argument('command', nargs=argparse.REMAINDER) options = parser.parse_args() command = options.command @@ -12,8 +13,11 @@ if not command: parser.error('a test command is required') -if options.disable_core_files: +if options.disable_core_files or options.max_virtual_memory is not None: resource.setrlimit(resource.RLIMIT_CORE, (0, 0)) +if options.max_virtual_memory is not None: + resource.setrlimit(resource.RLIMIT_AS, + (options.max_virtual_memory, options.max_virtual_memory + 1)) # Replace this process so test signals, timeouts and exit codes reach the runner. os.execvpe(command[0], command, os.environ) diff --git a/tools/test.py b/tools/test.py index ae2c82d12957..52e0a15745cd 100755 --- a/tools/test.py +++ b/tools/test.py @@ -889,19 +889,17 @@ def Execute(args, context, timeout=None, env=None, disable_core_files=False, # flags or environment variables defined via // Flags: and // Env: env_copy["NODE_SKIP_FLAG_CHECK"] = "true" - preexec_fn = None + resource_args = [] if disable_core_files and not utils.IsWindows(): - args = [sys.executable, join(dirname(__file__), 'test-resource-limits.py'), - '--disable-core-files', '--'] + args + resource_args.append('--disable-core-files') if max_virtual_memory is not None and utils.GuessOS() == 'linux': - def setMaxVirtualMemory(): - import resource - resource.setrlimit(resource.RLIMIT_CORE, (0,0)) - resource.setrlimit(resource.RLIMIT_AS, (max_virtual_memory,max_virtual_memory + 1)) + resource_args.extend(['--max-virtual-memory', str(max_virtual_memory)]) - preexec_fn = setMaxVirtualMemory + if resource_args: + args = [sys.executable, join(dirname(__file__), 'test-resource-limits.py')] + \ + resource_args + ['--'] + args (_process, exit_code, timed_out) = RunProcess( context, @@ -910,8 +908,7 @@ def setMaxVirtualMemory(): stdin = stdin, stdout = fd_out, stderr = fd_err, - env = env_copy, - preexec_fn = preexec_fn + env = env_copy ) os.close(fd_out) os.close(fd_err) From 3d6ac769a90e39bf0f71e04df3ddf936c24f3e98 Mon Sep 17 00:00:00 2001 From: Filip Skokan Date: Thu, 1 Oct 2026 14:26:41 +0200 Subject: [PATCH 05/10] test: isolate ports for concurrent benchmark tests Use OS-assigned listening ports for benchmark smoke tests and route clients to the bound addresses. Keep the existing port defaults for performance runs and let benchmark categories run in parallel. Select one valid compose smoke configuration to satisfy the existing single-configuration assertion. Signed-off-by: Filip Skokan Assisted-by: Codex --- benchmark/_http-benchmarkers.js | 4 +-- .../async_hooks/async-resource-vs-destroy.js | 1 + benchmark/async_hooks/http-server.js | 1 + benchmark/dgram/array-vs-concat.js | 6 +++-- benchmark/dgram/multi-buffer.js | 4 ++- benchmark/dgram/offset-length.js | 4 ++- benchmark/dgram/send-to-ip.js | 4 ++- benchmark/dgram/single-buffer.js | 4 ++- benchmark/diagnostics_channel/http.js | 1 + benchmark/http/cluster.js | 3 ++- benchmark/https/simple.js | 1 + benchmark/net/net-c2s-cork.js | 2 +- benchmark/net/net-c2s.js | 2 +- benchmark/net/net-pipe.js | 2 +- benchmark/net/net-s2c.js | 1 + benchmark/net/tcp-raw-c2s.js | 7 +++++- benchmark/net/tcp-raw-pipe.js | 7 +++++- benchmark/net/tcp-raw-s2c.js | 7 +++++- benchmark/streams/compose.js | 4 +-- benchmark/tls/secure-pair.js | 6 ++--- benchmark/tls/throughput-c2s.js | 2 +- benchmark/tls/throughput-s2c.js | 1 + benchmark/tls/tls-connect.js | 4 ++- test/README.md | 4 ++- test/benchmark/test-benchmark-http.js | 4 --- test/benchmark/test-benchmark-net.js | 4 --- test/benchmark/test-benchmark-tls.js | 4 --- test/benchmark/testcfg.py | 3 +-- test/common/benchmark.js | 2 +- test/parallel/test-benchmark-port.js | 25 +++++++++++++++++++ test/tools/test_test_configurations.py | 4 +-- 31 files changed, 88 insertions(+), 40 deletions(-) create mode 100644 test/parallel/test-benchmark-port.js diff --git a/benchmark/_http-benchmarkers.js b/benchmark/_http-benchmarkers.js index b9ef93c149f7..92f006c53923 100644 --- a/benchmark/_http-benchmarkers.js +++ b/benchmark/_http-benchmarkers.js @@ -7,8 +7,8 @@ const fs = require('fs'); const requirementsURL = 'https://github.com/nodejs/node/blob/HEAD/doc/contributing/writing-and-running-benchmarks.md#http-benchmark-requirements'; -// The port used by servers and wrk -exports.PORT = Number(process.env.PORT) || 12346; +// The port used by servers and wrk. Zero lets the OS select an available port. +exports.PORT = process.env.PORT === '0' ? 0 : Number(process.env.PORT) || 12346; class AutocannonBenchmarker { constructor() { diff --git a/benchmark/async_hooks/async-resource-vs-destroy.js b/benchmark/async_hooks/async-resource-vs-destroy.js index 3bd9dd56934e..f9b93cc39e30 100644 --- a/benchmark/async_hooks/async-resource-vs-destroy.js +++ b/benchmark/async_hooks/async-resource-vs-destroy.js @@ -176,6 +176,7 @@ function main({ type, asyncMethod, connections, duration, path }) { .on('listening', () => { bench.http({ + port: server.address().port, path, connections, duration, diff --git a/benchmark/async_hooks/http-server.js b/benchmark/async_hooks/http-server.js index da5f86f0e506..9b1b1106ffbe 100644 --- a/benchmark/async_hooks/http-server.js +++ b/benchmark/async_hooks/http-server.js @@ -32,6 +32,7 @@ function main({ asyncHooks, connections, duration }) { const path = '/buffer/4/4/normal/1'; bench.http({ + port: server.address().port, connections, path, duration, diff --git a/benchmark/dgram/array-vs-concat.js b/benchmark/dgram/array-vs-concat.js index c859771e7110..1cdc85b4b29f 100644 --- a/benchmark/dgram/array-vs-concat.js +++ b/benchmark/dgram/array-vs-concat.js @@ -25,6 +25,7 @@ function main({ dur, len, n, type, chunks }) { // Server let sent = 0; const socket = dgram.createSocket('udp4'); + let port; const onsend = type === 'concat' ? onsendConcat : onsendMulti; function onsendConcat() { @@ -33,7 +34,7 @@ function main({ dur, len, n, type, chunks }) { // that only perform synchronous I/O on nonblocking UDP sockets. setImmediate(() => { for (let i = 0; i < n; i++) { - socket.send(Buffer.concat(chunk), PORT, '127.0.0.1', onsend); + socket.send(Buffer.concat(chunk), port, '127.0.0.1', onsend); } }); } @@ -45,13 +46,14 @@ function main({ dur, len, n, type, chunks }) { // that only perform synchronous I/O on nonblocking UDP sockets. setImmediate(() => { for (let i = 0; i < n; i++) { - socket.send(chunk, PORT, '127.0.0.1', onsend); + socket.send(chunk, port, '127.0.0.1', onsend); } }); } } socket.on('listening', () => { + port = socket.address().port; bench.start(); onsend(); diff --git a/benchmark/dgram/multi-buffer.js b/benchmark/dgram/multi-buffer.js index 83d61daa36e1..9e578c52d3b4 100644 --- a/benchmark/dgram/multi-buffer.js +++ b/benchmark/dgram/multi-buffer.js @@ -24,6 +24,7 @@ function main({ dur, len, n, type, chunks }) { let sent = 0; let received = 0; const socket = dgram.createSocket('udp4'); + let port; function onsend() { if (sent++ % n === 0) { @@ -31,13 +32,14 @@ function main({ dur, len, n, type, chunks }) { // that only perform synchronous I/O on nonblocking UDP sockets. setImmediate(() => { for (let i = 0; i < n; i++) { - socket.send(chunk, PORT, '127.0.0.1', onsend); + socket.send(chunk, port, '127.0.0.1', onsend); } }); } } socket.on('listening', () => { + port = socket.address().port; bench.start(); onsend(); diff --git a/benchmark/dgram/offset-length.js b/benchmark/dgram/offset-length.js index 85381da8de78..2c6feabe206c 100644 --- a/benchmark/dgram/offset-length.js +++ b/benchmark/dgram/offset-length.js @@ -20,6 +20,7 @@ function main({ dur, len, n, type }) { let sent = 0; let received = 0; const socket = dgram.createSocket('udp4'); + let port; function onsend() { if (sent++ % n === 0) { @@ -27,13 +28,14 @@ function main({ dur, len, n, type }) { // that only perform synchronous I/O on nonblocking UDP sockets. setImmediate(() => { for (let i = 0; i < n; i++) { - socket.send(chunk, 0, chunk.length, PORT, '127.0.0.1', onsend); + socket.send(chunk, 0, chunk.length, port, '127.0.0.1', onsend); } }); } } socket.on('listening', () => { + port = socket.address().port; bench.start(); onsend(); diff --git a/benchmark/dgram/send-to-ip.js b/benchmark/dgram/send-to-ip.js index 8704462cf8d8..427487b49f11 100644 --- a/benchmark/dgram/send-to-ip.js +++ b/benchmark/dgram/send-to-ip.js @@ -18,18 +18,20 @@ function main({ dur, n }) { const chunk = Buffer.allocUnsafe(1); let sent = 0; const socket = dgram.createSocket('udp4'); + let port; function onsend() { if (sent++ % n === 0) { setImmediate(() => { for (let i = 0; i < n; i++) { - socket.send(chunk, PORT, '127.0.0.1', onsend); + socket.send(chunk, port, '127.0.0.1', onsend); } }); } } socket.on('listening', () => { + port = socket.address().port; bench.start(); onsend(); diff --git a/benchmark/dgram/single-buffer.js b/benchmark/dgram/single-buffer.js index 9ac41d15da8d..1031191a80c9 100644 --- a/benchmark/dgram/single-buffer.js +++ b/benchmark/dgram/single-buffer.js @@ -20,6 +20,7 @@ function main({ dur, len, n, type }) { let sent = 0; let received = 0; const socket = dgram.createSocket('udp4'); + let port; function onsend() { if (sent++ % n === 0) { @@ -27,13 +28,14 @@ function main({ dur, len, n, type }) { // that only perform synchronous I/O on nonblocking UDP sockets. setImmediate(() => { for (let i = 0; i < n; i++) { - socket.send(chunk, PORT, '127.0.0.1', onsend); + socket.send(chunk, port, '127.0.0.1', onsend); } }); } } socket.on('listening', () => { + port = socket.address().port; bench.start(); onsend(); diff --git a/benchmark/diagnostics_channel/http.js b/benchmark/diagnostics_channel/http.js index caf37a05a45c..fc304cdfa0e5 100644 --- a/benchmark/diagnostics_channel/http.js +++ b/benchmark/diagnostics_channel/http.js @@ -22,6 +22,7 @@ function main({ apm, connections, duration, type, len, chunks, chunkedEnc }) { .on('listening', () => { const path = `/${type}/${len}/${chunks}/normal/${chunkedEnc}`; bench.http({ + port: server.address().port, path, connections, duration, diff --git a/benchmark/http/cluster.js b/benchmark/http/cluster.js index 789bfc0e4dc7..2434b92e36bc 100644 --- a/benchmark/http/cluster.js +++ b/benchmark/http/cluster.js @@ -23,7 +23,7 @@ function main({ type, len, c, duration }) { const w1 = cluster.fork(); const w2 = cluster.fork(); - cluster.on('listening', () => { + cluster.on('listening', (worker, address) => { workers++; if (workers < 2) return; @@ -32,6 +32,7 @@ function main({ type, len, c, duration }) { const path = `/${type}/${len}`; bench.http({ + port: address.port, path: path, connections: c, duration, diff --git a/benchmark/https/simple.js b/benchmark/https/simple.js index 3b4af7caf631..5fdd14f68711 100644 --- a/benchmark/https/simple.js +++ b/benchmark/https/simple.js @@ -18,6 +18,7 @@ function main({ type, len, chunks, c, chunkedEnc, duration }) { const path = `/${type}/${len}/${chunks}/${chunkedEnc}`; bench.http({ + port: server.address().port, path, connections: c, scheme: 'https', diff --git a/benchmark/net/net-c2s-cork.js b/benchmark/net/net-c2s-cork.js index 9a1129218531..d39029b5b065 100644 --- a/benchmark/net/net-c2s-cork.js +++ b/benchmark/net/net-c2s-cork.js @@ -39,7 +39,7 @@ function main({ dur, len, type }) { }); server.listen(PORT, () => { - const socket = net.connect(PORT); + const socket = net.connect(server.address().port); socket.on('connect', () => { bench.start(); diff --git a/benchmark/net/net-c2s.js b/benchmark/net/net-c2s.js index 48e454c5815b..ff9822b1aa00 100644 --- a/benchmark/net/net-c2s.js +++ b/benchmark/net/net-c2s.js @@ -42,7 +42,7 @@ function main({ dur, len, type }) { }); server.listen(PORT, () => { - const socket = net.connect(PORT); + const socket = net.connect(server.address().port); socket.on('connect', () => { bench.start(); diff --git a/benchmark/net/net-pipe.js b/benchmark/net/net-pipe.js index dd9f2de497f5..73eb881eb859 100644 --- a/benchmark/net/net-pipe.js +++ b/benchmark/net/net-pipe.js @@ -42,7 +42,7 @@ function main({ dur, len, type }) { }); server.listen(PORT, () => { - const socket = net.connect(PORT); + const socket = net.connect(server.address().port); socket.on('connect', () => { bench.start(); diff --git a/benchmark/net/net-s2c.js b/benchmark/net/net-s2c.js index 1b9c91b8e780..185e720ccb96 100644 --- a/benchmark/net/net-s2c.js +++ b/benchmark/net/net-s2c.js @@ -76,6 +76,7 @@ function main({ dur, sendchunklen, type, recvbuflen, recvbufgenfn }) { }); server.listen(PORT, () => { + socketOpts.port = server.address().port; const socket = net.connect(socketOpts); socket.on('connect', () => { bench.start(); diff --git a/benchmark/net/tcp-raw-c2s.js b/benchmark/net/tcp-raw-c2s.js index 24e35906c925..b84d07fedafc 100644 --- a/benchmark/net/tcp-raw-c2s.js +++ b/benchmark/net/tcp-raw-c2s.js @@ -39,6 +39,11 @@ function main({ dur, len, type }) { if (err) fail(err, 'listen'); + const address = {}; + err = serverHandle.getsockname(address); + if (err) + fail(err, 'getsockname'); + serverHandle.onconnection = function(err, clientHandle) { if (err) fail(err, 'connect'); @@ -89,7 +94,7 @@ function main({ dur, len, type }) { const clientHandle = new TCP(TCPConstants.SOCKET); const connectReq = new TCPConnectWrap(); - const err = clientHandle.connect(connectReq, '127.0.0.1', PORT); + const err = clientHandle.connect(connectReq, '127.0.0.1', address.port); if (err) fail(err, 'connect'); diff --git a/benchmark/net/tcp-raw-pipe.js b/benchmark/net/tcp-raw-pipe.js index 36ddb822e04e..2cd306424ab7 100644 --- a/benchmark/net/tcp-raw-pipe.js +++ b/benchmark/net/tcp-raw-pipe.js @@ -45,6 +45,11 @@ function main({ dur, len, type }) { if (err) fail(err, 'listen'); + const address = {}; + err = serverHandle.getsockname(address); + if (err) + fail(err, 'getsockname'); + serverHandle.onconnection = function(err, clientHandle) { if (err) fail(err, 'connect'); @@ -93,7 +98,7 @@ function main({ dur, len, type }) { const connectReq = new TCPConnectWrap(); let bytes = 0; - err = clientHandle.connect(connectReq, '127.0.0.1', PORT); + err = clientHandle.connect(connectReq, '127.0.0.1', address.port); if (err) fail(err, 'connect'); diff --git a/benchmark/net/tcp-raw-s2c.js b/benchmark/net/tcp-raw-s2c.js index a847b8c28dcd..ed7eb3b66d3e 100644 --- a/benchmark/net/tcp-raw-s2c.js +++ b/benchmark/net/tcp-raw-s2c.js @@ -39,6 +39,11 @@ function main({ dur, len, type }) { if (err) fail(err, 'listen'); + const address = {}; + err = serverHandle.getsockname(address); + if (err) + fail(err, 'getsockname'); + serverHandle.onconnection = function(err, clientHandle) { if (err) fail(err, 'connect'); @@ -107,7 +112,7 @@ function main({ dur, len, type }) { function client(dur) { const clientHandle = new TCP(TCPConstants.SOCKET); const connectReq = new TCPConnectWrap(); - const err = clientHandle.connect(connectReq, '127.0.0.1', PORT); + const err = clientHandle.connect(connectReq, '127.0.0.1', address.port); if (err) fail(err, 'connect'); diff --git a/benchmark/streams/compose.js b/benchmark/streams/compose.js index 283ad8b7e30b..461c8282eebd 100644 --- a/benchmark/streams/compose.js +++ b/benchmark/streams/compose.js @@ -18,8 +18,8 @@ const bench = common.createBenchmark(main, { return type === 'creation' ? n === 1e3 : n === 1; }, test: { - n: [1, 1e3], - type: ['creation', 'throughput'], + n: 1, + type: 'throughput', }, }); diff --git a/benchmark/tls/secure-pair.js b/benchmark/tls/secure-pair.js index a253bbf02607..1fce51707797 100644 --- a/benchmark/tls/secure-pair.js +++ b/benchmark/tls/secure-pair.js @@ -12,7 +12,7 @@ const fixtures = require('../../test/common/fixtures'); const tls = require('tls'); const net = require('net'); -const REDIRECT_PORT = 28347; +const REDIRECT_PORT = common.PORT === 0 ? 0 : 28347; function main({ dur, size, securing }) { const chunk = Buffer.alloc(size, 'b'); @@ -33,7 +33,7 @@ function main({ dur, size, securing }) { const proxy = net.createServer(onProxyConnection); proxy.listen(common.PORT, () => { const clientOptions = { - port: common.PORT, + port: proxy.address().port, ca: options.ca, key: options.key, cert: options.cert, @@ -66,7 +66,7 @@ function main({ dur, size, securing }) { }); function onProxyConnection(conn) { - const client = net.connect(REDIRECT_PORT, () => { + const client = net.connect(server.address().port, () => { switch (securing) { case 'TLSSocket': secureTLSSocket(conn, client); diff --git a/benchmark/tls/throughput-c2s.js b/benchmark/tls/throughput-c2s.js index bf71f92aecbc..c3d8d4390311 100644 --- a/benchmark/tls/throughput-c2s.js +++ b/benchmark/tls/throughput-c2s.js @@ -40,7 +40,7 @@ function main({ dur, type, size }) { const server = tls.createServer(options, onConnection); let conn; server.listen(common.PORT, () => { - const opt = { port: common.PORT, rejectUnauthorized: false }; + const opt = { port: server.address().port, rejectUnauthorized: false }; conn = tls.connect(opt, () => { setTimeout(done, dur * 1000); bench.start(); diff --git a/benchmark/tls/throughput-s2c.js b/benchmark/tls/throughput-s2c.js index 7fb93c304b20..e2159c8f68d7 100644 --- a/benchmark/tls/throughput-s2c.js +++ b/benchmark/tls/throughput-s2c.js @@ -86,6 +86,7 @@ function main({ dur, type, sendchunklen, recvbuflen, recvbufgenfn }) { let conn; server.listen(common.PORT, () => { + socketOpts.port = server.address().port; conn = tls.connect(socketOpts, () => { setTimeout(done, dur * 1000); bench.start(); diff --git a/benchmark/tls/tls-connect.js b/benchmark/tls/tls-connect.js index b398bdb3c3d1..3ae2a8fe6527 100644 --- a/benchmark/tls/tls-connect.js +++ b/benchmark/tls/tls-connect.js @@ -12,6 +12,7 @@ let clientConn = 0; let serverConn = 0; let dur; let concurrency; +let port; let running = true; function main(conf) { @@ -30,6 +31,7 @@ function main(conf) { } function onListening() { + port = this.address().port; setTimeout(done, dur * 1000); bench.start(); for (let i = 0; i < concurrency; i++) @@ -42,7 +44,7 @@ function onConnection(conn) { function makeConnection() { const options = { - port: common.PORT, + port, rejectUnauthorized: false, }; const conn = tls.connect(options, () => { diff --git a/test/README.md b/test/README.md index 9b2f2ca7e93d..d7451af7beaa 100644 --- a/test/README.md +++ b/test/README.md @@ -40,10 +40,12 @@ For the tests to run on Windows, be sure to clone Node.js source code with the Tests run in parallel by default. Suite configurations opt out with `SerialTestConfiguration`, or `SerialAddonTestConfiguration` for addon layouts. -The `addons`, `benchmark`, `internet`, `js-native-api`, `known_issues`, `node-api`, +The `addons`, `internet`, `js-native-api`, `known_issues`, `node-api`, and `pummel` suites explicitly run serially. WPT's group settings and SEA's disk-space guard retain their existing scheduling. +Benchmark smoke tests use TCP/UDP ports selected by the OS. + The test runner finishes parallel tests first, followed by serial suites, then `sequential` tests. In `sequential`, different subsystems can run concurrently up to the worker count selected with `-j`. The subsystem is the first filename diff --git a/test/benchmark/test-benchmark-http.js b/test/benchmark/test-benchmark-http.js index a3d92c7e987f..9aac184d5332 100644 --- a/test/benchmark/test-benchmark-http.js +++ b/test/benchmark/test-benchmark-http.js @@ -5,10 +5,6 @@ const common = require('../common'); if (!common.enoughTestMem) common.skip('Insufficient memory for HTTP benchmark test'); -// Because the http benchmarks use hardcoded ports, this should be in sequential -// rather than parallel to make sure it does not conflict with tests that choose -// random available ports. - const runBenchmark = require('../common/benchmark'); runBenchmark('http', { NODEJS_BENCHMARK_ZERO_ALLOWED: 1 }); diff --git a/test/benchmark/test-benchmark-net.js b/test/benchmark/test-benchmark-net.js index df8ea8011693..f791b8a5f018 100644 --- a/test/benchmark/test-benchmark-net.js +++ b/test/benchmark/test-benchmark-net.js @@ -2,10 +2,6 @@ require('../common'); -// Because the net benchmarks use hardcoded ports, this should be in sequential -// rather than parallel to make sure it does not conflict with tests that choose -// random available ports. - const runBenchmark = require('../common/benchmark'); runBenchmark('net', { NODEJS_BENCHMARK_ZERO_ALLOWED: 1 }); diff --git a/test/benchmark/test-benchmark-tls.js b/test/benchmark/test-benchmark-tls.js index c9a87c15770d..5769516b2506 100644 --- a/test/benchmark/test-benchmark-tls.js +++ b/test/benchmark/test-benchmark-tls.js @@ -8,10 +8,6 @@ if (!common.hasCrypto) if (!common.enoughTestMem) common.skip('Insufficient memory for TLS benchmark test'); -// Because the TLS benchmarks use hardcoded ports, this should be in sequential -// rather than parallel to make sure it does not conflict with tests that choose -// random available ports. - const runBenchmark = require('../common/benchmark'); runBenchmark('tls', { NODEJS_BENCHMARK_ZERO_ALLOWED: 1 }); diff --git a/test/benchmark/testcfg.py b/test/benchmark/testcfg.py index cf364de83f0b..2c2929f610b8 100644 --- a/test/benchmark/testcfg.py +++ b/test/benchmark/testcfg.py @@ -3,5 +3,4 @@ import testpy def GetConfiguration(context, root): - # TODO: Isolate fixed TCP/UDP ports before parallelizing benchmark tests. - return testpy.SerialTestConfiguration(context, root, 'benchmark') + return testpy.SimpleTestConfiguration(context, root, 'benchmark') diff --git a/test/common/benchmark.js b/test/common/benchmark.js index fe8ac8d7c0a4..8ad0ae695aa5 100644 --- a/test/common/benchmark.js +++ b/test/common/benchmark.js @@ -11,7 +11,7 @@ function runBenchmark(name, env) { argv.push(name); - const mergedEnv = { ...process.env, ...env }; + const mergedEnv = { ...process.env, ...env, PORT: '0' }; const child = fork(runjs, argv, { env: mergedEnv, diff --git a/test/parallel/test-benchmark-port.js b/test/parallel/test-benchmark-port.js new file mode 100644 index 000000000000..3c8ea231f491 --- /dev/null +++ b/test/parallel/test-benchmark-port.js @@ -0,0 +1,25 @@ +'use strict'; + +require('../common'); +const assert = require('assert'); +const { spawnSync } = require('child_process'); +const path = require('path'); + +const benchmarkers = path.resolve(__dirname, '../../benchmark/_http-benchmarkers.js'); +const script = `console.log(require(${JSON.stringify(benchmarkers)}).PORT)`; + +for (const [port, expected] of [ + [undefined, 12346], + ['', 12346], + ['0', 0], + ['31337', 31337], +]) { + const env = { ...process.env }; + delete env.PORT; + if (port !== undefined) + env.PORT = port; + const child = spawnSync(process.execPath, ['-e', script], { env, encoding: 'utf8' }); + assert.ifError(child.error); + assert.strictEqual(child.status, 0, child.stderr); + assert.strictEqual(child.stdout, `${expected}\n`); +} diff --git a/test/tools/test_test_configurations.py b/test/tools/test_test_configurations.py index 5d92883d7e41..401cb35556cc 100644 --- a/test/tools/test_test_configurations.py +++ b/test/tools/test_test_configurations.py @@ -57,7 +57,7 @@ def test_generic_and_serial_addon_configurations_preserve_discovery(self): self.assertTrue(all(case.parallel == parallel for case in cases)) def test_named_suites_explicitly_remain_serial(self): - for suite in ['pummel', 'benchmark', 'known_issues', 'internet', + for suite in ['pummel', 'known_issues', 'internet', 'sequential', 'addons', 'js-native-api', 'node-api']: with self.subTest(suite=suite): @@ -67,7 +67,7 @@ def test_named_suites_explicitly_remain_serial(self): self.assertTrue(all(not case.parallel for case in cases)) def test_other_suites_use_parallel_default(self): - for suite in ['parallel', 'abort', 'async-hooks', 'client-proxy', 'doctool', + for suite in ['parallel', 'abort', 'async-hooks', 'benchmark', 'client-proxy', 'doctool', 'embedding', 'es-module', 'ffi', 'module-hooks', 'report', 'sqlite', 'test426', 'test-runner', 'tick-processor', 'trace_events', 'v8-updates', 'wasi', 'wasm-allocation']: From 2f914c3717a17100d8d9f857d28e6a3c6b8de495 Mon Sep 17 00:00:00 2001 From: Filip Skokan Date: Thu, 1 Oct 2026 14:27:46 +0200 Subject: [PATCH 06/10] test: isolate internet test listening ports Bind UDP senders to OS-assigned ports and distribute those ports to child listeners over IPC. Reserve each Windows multicast port for its test and use an ephemeral inspector port, then run the internet suite in parallel. Signed-off-by: Filip Skokan Assisted-by: Codex --- test/README.md | 3 +- .../test-dgram-broadcast-multi-process.js | 13 +++-- test/internet/test-dgram-connect.js | 7 +-- .../test-dgram-multicast-multi-process.js | 16 +++++-- .../test-dgram-multicast-set-interface-lo.js | 48 ++++++++++++++----- .../test-dgram-multicast-ssm-multi-process.js | 16 +++++-- ...est-dgram-multicast-ssmv6-multi-process.js | 16 +++++-- test/internet/test-inspector-help-page.js | 2 +- test/internet/testcfg.py | 3 +- test/tools/test_test_configurations.py | 4 +- 10 files changed, 88 insertions(+), 40 deletions(-) diff --git a/test/README.md b/test/README.md index d7451af7beaa..f7d51cff3b5f 100644 --- a/test/README.md +++ b/test/README.md @@ -40,11 +40,12 @@ For the tests to run on Windows, be sure to clone Node.js source code with the Tests run in parallel by default. Suite configurations opt out with `SerialTestConfiguration`, or `SerialAddonTestConfiguration` for addon layouts. -The `addons`, `internet`, `js-native-api`, `known_issues`, `node-api`, +The `addons`, `js-native-api`, `known_issues`, `node-api`, and `pummel` suites explicitly run serially. WPT's group settings and SEA's disk-space guard retain their existing scheduling. Benchmark smoke tests use TCP/UDP ports selected by the OS. +Internet tests also allocate their listening ports dynamically. The test runner finishes parallel tests first, followed by serial suites, then `sequential` tests. In `sequential`, different subsystems can run concurrently diff --git a/test/internet/test-dgram-broadcast-multi-process.js b/test/internet/test-dgram-broadcast-multi-process.js index aa6ef56e9f4d..533f3a40492d 100644 --- a/test/internet/test-dgram-broadcast-multi-process.js +++ b/test/internet/test-dgram-broadcast-multi-process.js @@ -63,6 +63,7 @@ if (process.argv[2] !== 'child') { let i = 0; let done = 0; let timer = null; + let port; // Exit the test if it doesn't succeed within TIMEOUT timer = setTimeout(() => { @@ -171,9 +172,13 @@ if (process.argv[2] !== 'child') { // Bind the address explicitly for sending // INADDR_BROADCAST to only one interface - sendSocket.bind(common.PORT, bindAddress); + sendSocket.bind(0, bindAddress); sendSocket.on('listening', () => { sendSocket.setBroadcast(true); + port = sendSocket.address().port; + for (const worker of Object.values(workers)) { + worker.send(port); + } }); sendSocket.on('close', () => { @@ -194,12 +199,12 @@ if (process.argv[2] !== 'child') { buf, 0, buf.length, - common.PORT, + port, LOCAL_BROADCAST_HOST, common.mustSucceed(() => { console.error('[PARENT] sent %s to %s:%s', util.inspect(buf.toString()), - LOCAL_BROADCAST_HOST, common.PORT); + LOCAL_BROADCAST_HOST, port); process.nextTick(sendSocket.sendNext); }), @@ -248,5 +253,5 @@ if (process.argv[2] === 'child') { listenSocket.on('listening', () => { process.send({ listening: true }); }); - listenSocket.bind(common.PORT); + process.once('message', common.mustCall((port) => listenSocket.bind(port))); } diff --git a/test/internet/test-dgram-connect.js b/test/internet/test-dgram-connect.js index 47a12c789092..c3a973b3405d 100644 --- a/test/internet/test-dgram-connect.js +++ b/test/internet/test-dgram-connect.js @@ -4,16 +4,17 @@ const common = require('../common'); const { addresses } = require('../common/internet'); const assert = require('assert'); const dgram = require('dgram'); +const port = 12345; const client = dgram.createSocket('udp4'); -client.connect(common.PORT, addresses.INVALID_HOST, common.mustCall((err) => { +client.connect(port, addresses.INVALID_HOST, common.mustCall((err) => { assert.ok(err.code === 'ENOTFOUND' || err.code === 'EAI_AGAIN'); client.once('error', common.mustCall((err) => { assert.ok(err.code === 'ENOTFOUND' || err.code === 'EAI_AGAIN'); client.once('connect', common.mustCall(() => client.close())); - client.connect(common.PORT); + client.connect(port); })); - client.connect(common.PORT, addresses.INVALID_HOST); + client.connect(port, addresses.INVALID_HOST); })); diff --git a/test/internet/test-dgram-multicast-multi-process.js b/test/internet/test-dgram-multicast-multi-process.js index 20051dbb07ec..370fd4e3c06f 100644 --- a/test/internet/test-dgram-multicast-multi-process.js +++ b/test/internet/test-dgram-multicast-multi-process.js @@ -39,7 +39,7 @@ const messages = [ ]; const workers = {}; const listeners = 3; -let listening, sendSocket, done, timer, dead; +let listening, sendSocket, done, timer, dead, port; function launchChildProcess() { @@ -151,7 +151,9 @@ if (process.argv[2] !== 'child') { launchChildProcess(x); } - sendSocket = dgram.createSocket('udp4'); + // Share an ephemeral port with the receivers in this test. + sendSocket = dgram.createSocket({ type: 'udp4', reuseAddr: true }); + sendSocket.bind(0); // The socket is actually created async now. sendSocket.on('listening', function() { @@ -160,6 +162,10 @@ if (process.argv[2] !== 'child') { sendSocket.setMulticastTTL(1); sendSocket.setMulticastLoopback(true); sendSocket.setMulticastInterface(LOCAL_HOST_IFADDR); + port = sendSocket.address().port; + for (const worker of Object.values(workers)) { + worker.send(port); + } }); sendSocket.on('close', function() { @@ -180,12 +186,12 @@ if (process.argv[2] !== 'child') { buf, 0, buf.length, - common.PORT, + port, LOCAL_BROADCAST_HOST, common.mustSucceed(() => { console.error('[PARENT] sent "%s" to %s:%s', buf.toString(), - LOCAL_BROADCAST_HOST, common.PORT); + LOCAL_BROADCAST_HOST, port); process.nextTick(sendSocket.sendNext); }), ); @@ -230,5 +236,5 @@ if (process.argv[2] === 'child') { process.send({ listening: true }); }); - listenSocket.bind(common.PORT); + process.once('message', common.mustCall((port) => listenSocket.bind(port))); } diff --git a/test/internet/test-dgram-multicast-set-interface-lo.js b/test/internet/test-dgram-multicast-set-interface-lo.js index 4850ffca9409..fc9efbb4464c 100644 --- a/test/internet/test-dgram-multicast-set-interface-lo.js +++ b/test/internet/test-dgram-multicast-set-interface-lo.js @@ -31,11 +31,7 @@ const LOOPBACK = { IPv4: '127.0.0.1', IPv6: '::1' }; const ANY = { IPv4: '0.0.0.0', IPv6: '::' }; const FAM = 'IPv4'; -// Windows won't bind on multicasts so its filtering is by port. const PORTS = {}; -for (let i = 0; i < MULTICASTS[FAM].length; i++) { - PORTS[MULTICASTS[FAM][i]] = common.PORT + (common.isWindows ? i : 0); -} const UDP = { IPv4: 'udp4', IPv6: 'udp6' }; @@ -116,6 +112,7 @@ if (process.argv[2] !== 'child') { messagesNeeded.length, NOW]); workers[worker.pid] = worker; + worker.multicast = MULTICAST; worker.messagesReceived = []; worker.messagesNeeded = messagesNeeded; @@ -204,9 +201,32 @@ if (process.argv[2] !== 'child') { reuseAddr: true, }); - // Don't bind the address explicitly when sending and start with - // the OSes default multicast interface selection. - sendSocket.bind(common.PORT, ANY[FAM]); + // Reserve the shared port before allowing the children to bind it. + // Windows won't bind on multicasts so its filtering is by port. + const sendSockets = [sendSocket]; + if (common.isWindows) { + for (let i = 1; i < MULTICASTS[FAM].length; i++) { + sendSockets.push(dgram.createSocket({ type: UDP[FAM], reuseAddr: true })); + } + } + for (const [i, socket] of sendSockets.entries()) { + socket.on('listening', () => { + if (common.isWindows) { + PORTS[MULTICASTS[FAM][i]] = socket.address().port; + } else { + for (const multicast of MULTICASTS[FAM]) { + PORTS[multicast] = socket.address().port; + } + } + if (Object.keys(PORTS).length === MULTICASTS[FAM].length) { + for (const worker of Object.values(workers)) { + worker.send(PORTS[worker.multicast]); + } + } + }); + // Start with the OS's default multicast interface selection. + socket.bind(0, ANY[FAM]); + } sendSocket.on('listening', () => { console.error(`outgoing iface ${interfaceAddress}`); }); @@ -219,7 +239,9 @@ if (process.argv[2] !== 'child') { const msg = messages[i++]; if (!msg) { - sendSocket.close(); + for (const socket of sendSockets) { + socket.close(); + } return; } console.error(TMPL(NOW, msg.tail)); @@ -285,8 +307,10 @@ if (process.argv[2] === 'child') { process.send({ listening: true }); }); - if (common.isWindows) - listenSocket.bind(PORTS[MULTICAST], ANY[FAM]); - else - listenSocket.bind(common.PORT, MULTICAST); + process.once('message', common.mustCall((port) => { + if (common.isWindows) + listenSocket.bind(port, ANY[FAM]); + else + listenSocket.bind(port, MULTICAST); + })); } diff --git a/test/internet/test-dgram-multicast-ssm-multi-process.js b/test/internet/test-dgram-multicast-ssm-multi-process.js index e363b6234fae..69b0dd8b4038 100644 --- a/test/internet/test-dgram-multicast-ssm-multi-process.js +++ b/test/internet/test-dgram-multicast-ssm-multi-process.js @@ -18,7 +18,7 @@ const messages = [ ]; const workers = {}; const listeners = 3; -let listening, sendSocket, done, timer, dead; +let listening, sendSocket, done, timer, dead, port; let sourceAddress = null; @@ -145,7 +145,9 @@ if (process.argv[2] !== 'child') { launchChildProcess(x); } - sendSocket = dgram.createSocket('udp4'); + // Share an ephemeral port with the receivers in this test. + sendSocket = dgram.createSocket({ type: 'udp4', reuseAddr: true }); + sendSocket.bind(0); // The socket is actually created async now. sendSocket.on('listening', function() { @@ -154,6 +156,10 @@ if (process.argv[2] !== 'child') { sendSocket.setMulticastTTL(1); sendSocket.setMulticastLoopback(true); sendSocket.addSourceSpecificMembership(sourceAddress, GROUP_ADDRESS); + port = sendSocket.address().port; + for (const worker of Object.values(workers)) { + worker.send(port); + } }); sendSocket.on('close', function() { @@ -174,12 +180,12 @@ if (process.argv[2] !== 'child') { buf, 0, buf.length, - common.PORT, + port, GROUP_ADDRESS, common.mustSucceed((err) => { console.error('[PARENT] sent "%s" to %s:%s', buf.toString(), - GROUP_ADDRESS, common.PORT); + GROUP_ADDRESS, port); process.nextTick(sendSocket.sendNext); }), ); @@ -226,5 +232,5 @@ if (process.argv[2] === 'child') { process.send({ listening: true }); }); - listenSocket.bind(common.PORT); + process.once('message', common.mustCall((port) => listenSocket.bind(port))); } diff --git a/test/internet/test-dgram-multicast-ssmv6-multi-process.js b/test/internet/test-dgram-multicast-ssmv6-multi-process.js index fda5fa5fab06..81380cee65ac 100644 --- a/test/internet/test-dgram-multicast-ssmv6-multi-process.js +++ b/test/internet/test-dgram-multicast-ssmv6-multi-process.js @@ -18,7 +18,7 @@ const messages = [ ]; const workers = {}; const listeners = 3; -let listening, sendSocket, done, timer, dead; +let listening, sendSocket, done, timer, dead, port; let sourceAddress = null; @@ -145,7 +145,9 @@ if (process.argv[2] !== 'child') { launchChildProcess(x); } - sendSocket = dgram.createSocket('udp6'); + // Share an ephemeral port with the receivers in this test. + sendSocket = dgram.createSocket({ type: 'udp6', reuseAddr: true }); + sendSocket.bind(0); // The socket is actually created async now. sendSocket.on('listening', function() { @@ -154,6 +156,10 @@ if (process.argv[2] !== 'child') { sendSocket.setMulticastTTL(1); sendSocket.setMulticastLoopback(true); sendSocket.addSourceSpecificMembership(sourceAddress, GROUP_ADDRESS); + port = sendSocket.address().port; + for (const worker of Object.values(workers)) { + worker.send(port); + } }); sendSocket.on('close', function() { @@ -174,12 +180,12 @@ if (process.argv[2] !== 'child') { buf, 0, buf.length, - common.PORT, + port, GROUP_ADDRESS, common.mustSucceed(() => { console.error('[PARENT] sent "%s" to %s:%s', buf.toString(), - GROUP_ADDRESS, common.PORT); + GROUP_ADDRESS, port); process.nextTick(sendSocket.sendNext); }), ); @@ -226,5 +232,5 @@ if (process.argv[2] === 'child') { process.send({ listening: true }); }); - listenSocket.bind(common.PORT); + process.once('message', common.mustCall((port) => listenSocket.bind(port))); } diff --git a/test/internet/test-inspector-help-page.js b/test/internet/test-inspector-help-page.js index ae06242405d5..41e14816c297 100644 --- a/test/internet/test-inspector-help-page.js +++ b/test/internet/test-inspector-help-page.js @@ -9,7 +9,7 @@ if (!common.hasCrypto) const assert = require('assert'); const https = require('https'); const { spawnSync } = require('child_process'); -const child = spawnSync(process.execPath, ['--inspect', '-e', '""']); +const child = spawnSync(process.execPath, ['--inspect=0', '-e', '""']); const stderr = child.stderr.toString(); const helpUrl = stderr.match(/For help, see: (.+)/)[1]; diff --git a/test/internet/testcfg.py b/test/internet/testcfg.py index 7c68a8573976..73e70e340000 100644 --- a/test/internet/testcfg.py +++ b/test/internet/testcfg.py @@ -3,5 +3,4 @@ import testpy def GetConfiguration(context, root): - # TODO: Isolate shared listening ports before allowing concurrent tests. - return testpy.SerialTestConfiguration(context, root, 'internet') + return testpy.SimpleTestConfiguration(context, root, 'internet') diff --git a/test/tools/test_test_configurations.py b/test/tools/test_test_configurations.py index 401cb35556cc..0d4532c34f29 100644 --- a/test/tools/test_test_configurations.py +++ b/test/tools/test_test_configurations.py @@ -57,7 +57,7 @@ def test_generic_and_serial_addon_configurations_preserve_discovery(self): self.assertTrue(all(case.parallel == parallel for case in cases)) def test_named_suites_explicitly_remain_serial(self): - for suite in ['pummel', 'known_issues', 'internet', + for suite in ['pummel', 'known_issues', 'sequential', 'addons', 'js-native-api', 'node-api']: with self.subTest(suite=suite): @@ -68,7 +68,7 @@ def test_named_suites_explicitly_remain_serial(self): def test_other_suites_use_parallel_default(self): for suite in ['parallel', 'abort', 'async-hooks', 'benchmark', 'client-proxy', 'doctool', - 'embedding', 'es-module', 'ffi', 'module-hooks', 'report', + 'embedding', 'es-module', 'ffi', 'internet', 'module-hooks', 'report', 'sqlite', 'test426', 'test-runner', 'tick-processor', 'trace_events', 'v8-updates', 'wasi', 'wasm-allocation']: with self.subTest(suite=suite): From 0cda872cea8240ba978c8cd0c82981700f96df0f Mon Sep 17 00:00:00 2001 From: Filip Skokan Date: Thu, 1 Oct 2026 14:27:51 +0200 Subject: [PATCH 07/10] test: schedule isolated WPT groups in parallel Separate managed-process scheduling from standalone worker concurrency. Allow web-locks and webstorage groups to overlap because each process has its own LockManager and storage directory. Keep standalone workers and timing-sensitive timer groups serial. Signed-off-by: Filip Skokan Assisted-by: Codex --- test/README.md | 5 +++-- test/common/wpt.js | 3 ++- test/parallel/test-common-wpt-runner.js | 4 +++- test/wpt/test-timers.js | 3 ++- test/wpt/test-web-locks.js | 2 +- test/wpt/test-webstorage.js | 1 + test/wpt/testcfg.py | 1 - 7 files changed, 12 insertions(+), 7 deletions(-) diff --git a/test/README.md b/test/README.md index f7d51cff3b5f..42a1cbe3ae1e 100644 --- a/test/README.md +++ b/test/README.md @@ -41,8 +41,9 @@ For the tests to run on Windows, be sure to clone Node.js source code with the Tests run in parallel by default. Suite configurations opt out with `SerialTestConfiguration`, or `SerialAddonTestConfiguration` for addon layouts. The `addons`, `js-native-api`, `known_issues`, `node-api`, -and `pummel` suites explicitly run serially. WPT's group settings and SEA's -disk-space guard retain their existing scheduling. +and `pummel` suites explicitly run serially. WPT timer groups also run serially; +managed web-locks and webstorage groups have isolated processes and run in +parallel. SEA's disk-space guard uses the requested worker count. Benchmark smoke tests use TCP/UDP ports selected by the OS. Internet tests also allocate their listening ports dynamically. diff --git a/test/common/wpt.js b/test/common/wpt.js index 29180fdeb358..f71a7eeabaa3 100644 --- a/test/common/wpt.js +++ b/test/common/wpt.js @@ -835,7 +835,8 @@ class WPTRunner { } } this.isListing = this.managed?.mode === 'list'; - this.serial = options.concurrency === 1; + // Managed groups use separate processes, independent of worker concurrency. + this.serial = options.serial === true; if (this.managed?.mode === 'run') concurrency = 1; // RISC-V has very limited virtual address space in the currently common diff --git a/test/parallel/test-common-wpt-runner.js b/test/parallel/test-common-wpt-runner.js index ea41caa6cefc..3fd4561b46c6 100644 --- a/test/parallel/test-common-wpt-runner.js +++ b/test/parallel/test-common-wpt-runner.js @@ -103,13 +103,15 @@ function main() { assert.match(skippedOutput, /1\.\.0 # SKIP/); assert.doesNotMatch(skippedOutput, /\[PASS\]/); + assert.strictEqual(discover('web-locks').serial, false); + if (common.hasSQLite) { const root = path.join(tmpdir.path, 'discovery'); const directory = path.join(root, '.tmp.probe'); fs.mkdirSync(directory, { recursive: true }); const sentinel = path.join(directory, 'sentinel'); fs.writeFileSync(sentinel, 'preserved'); - assert.strictEqual(discover('webstorage', { NODE_TEST_DIR: root, TEST_SERIAL_ID: 'probe' }).serial, true); + assert.strictEqual(discover('webstorage', { NODE_TEST_DIR: root, TEST_SERIAL_ID: 'probe' }).serial, false); assert.strictEqual(fs.readFileSync(sentinel, 'utf8'), 'preserved'); } diff --git a/test/wpt/test-timers.js b/test/wpt/test-timers.js index ba84de258324..5ed9a1c44b6a 100644 --- a/test/wpt/test-timers.js +++ b/test/wpt/test-timers.js @@ -4,7 +4,8 @@ const assert = require('assert'); const { basename } = require('path'); const { WPTRunner } = require('../common/wpt'); -const runner = new WPTRunner('html/webappapis/timers', { concurrency: 1 }); +// Keep timing-sensitive tests out of the parallel runner phase. +const runner = new WPTRunner('html/webappapis/timers', { concurrency: 1, serial: true }); runner.setScriptModifier((script) => { if (!['type-long-settimeout.any.js', 'type-long-setinterval.any.js'] diff --git a/test/wpt/test-web-locks.js b/test/wpt/test-web-locks.js index f7080a23757d..7feada5a464d 100644 --- a/test/wpt/test-web-locks.js +++ b/test/wpt/test-web-locks.js @@ -2,7 +2,7 @@ const { WPTRunner } = require('../common/wpt'); -// Run serially to avoid cross-test interference on the shared LockManager. +// Standalone worker threads share a LockManager; managed groups have their own process. const runner = new WPTRunner('web-locks', { concurrency: 1 }); runner.pretendGlobalThisAs('Window'); diff --git a/test/wpt/test-webstorage.js b/test/wpt/test-webstorage.js index b0cc3d6bf13d..0da2684bddb3 100644 --- a/test/wpt/test-webstorage.js +++ b/test/wpt/test-webstorage.js @@ -4,6 +4,7 @@ skipIfSQLiteMissing(); const tmpdir = require('../common/tmpdir'); const { WPTRunner } = require('../common/wpt'); const { join } = require('node:path'); +// Standalone workers share this file; managed groups use separate test directories. const runner = new WPTRunner('webstorage', { concurrency: 1 }); if (!runner.isListing) tmpdir.refresh(); diff --git a/test/wpt/testcfg.py b/test/wpt/testcfg.py index 80ed22f7d3f3..f11ef1f34cbb 100644 --- a/test/wpt/testcfg.py +++ b/test/wpt/testcfg.py @@ -17,7 +17,6 @@ def __init__(self, path, file, arch, mode, context, config, group, serial): super(WPTTestCase, self).__init__( path, file, arch, mode, context, config, config.additional_flags) self.group = group - # TODO: Audit serial groups for process isolation and timer sensitivity. self.parallel = not serial def GetName(self): From 7c0570a13ff7f4ca63cd69b9996a3ea5d12dacb7 Mon Sep 17 00:00:00 2001 From: Filip Skokan Date: Thu, 1 Oct 2026 15:38:24 +0200 Subject: [PATCH 08/10] test: tolerate invalid UTF-8 in test output Malformed bytes in stdout or stderr currently raise UnicodeDecodeError and stop the test runner. Decode with replacement characters so the runner can report the test's exit status and diagnostics. Signed-off-by: Filip Skokan Assisted-by: Codex --- test/tools/test_test_runner.py | 25 +++++++++++++++++++++++++ tools/test.py | 4 ++-- 2 files changed, 27 insertions(+), 2 deletions(-) diff --git a/test/tools/test_test_runner.py b/test/tools/test_test_runner.py index 29ffb93e8775..9bfd3ba172eb 100644 --- a/test/tools/test_test_runner.py +++ b/test/tools/test_test_runner.py @@ -4,6 +4,7 @@ import sys import threading import unittest +from types import SimpleNamespace from unittest import mock ROOT = os.path.abspath(os.path.join(os.path.dirname(__file__), '..', '..')) @@ -58,6 +59,30 @@ def HasRun(self, output): event.set() +class ExecuteTest(unittest.TestCase): + def test_invalid_utf8_output_preserves_test_result(self): + context = SimpleNamespace(verbose=False, suppress_dialogs=False, + abort_on_timeout=False) + code = '''import sys +sys.stdout.buffer.write(b'valid: \\xc3\\xa9\\ninvalid: \\xe8!\\n') +sys.stderr.buffer.write(b'valid: \\xe2\\x82\\xac\\ninvalid: \\xff!\\n') +sys.exit(7) +''' + output = runner.Execute([sys.executable, '-c', code], context, timeout=5) + self.assertEqual(output.stdout, 'valid: é\ninvalid: \ufffd!\n') + self.assertEqual(output.stderr, 'valid: €\ninvalid: \ufffd!\n') + self.assertEqual(output.exit_code, 7) + self.assertFalse(output.timed_out) + + case = SchedulerCase('test-net-invalid-utf8') + failure = runner.TestOutput(case, ['node', 'test-net-invalid-utf8'], output, False) + self.assertTrue(failure.UnexpectedOutput()) + report = SchedulerProgress([case]).GetFailureOutput(failure) + self.assertIn(output.stdout.strip(), report) + self.assertIn(output.stderr.strip(), report) + report.encode('utf8') + + class SchedulerTest(unittest.TestCase): def wait_for(self, event): self.assertTrue(event.wait(WAIT_TIMEOUT), 'test runner did not make progress') diff --git a/tools/test.py b/tools/test.py index 52e0a15745cd..f8bfbb8f6cf2 100755 --- a/tools/test.py +++ b/tools/test.py @@ -912,9 +912,9 @@ def Execute(args, context, timeout=None, env=None, disable_core_files=False, ) os.close(fd_out) os.close(fd_err) - with open(outname, encoding='utf8') as output_file: + with open(outname, encoding='utf8', errors='replace') as output_file: output = output_file.read() - with open(errname, encoding='utf8') as error_file: + with open(errname, encoding='utf8', errors='replace') as error_file: errors = error_file.read() CheckedUnlink(outname) CheckedUnlink(errname) From 127744ca8c503f76f7afb82d58b95e7c1b8c9885 Mon Sep 17 00:00:00 2001 From: Filip Skokan Date: Thu, 1 Oct 2026 15:47:06 +0200 Subject: [PATCH 09/10] test: isolate test-pipe listening ports The pipe subsystem can overlap with other tests using common.PORT. Let the OS select both listening ports and connect clients to the bound addresses, preserving the existing transfer assertions. Signed-off-by: Filip Skokan Assisted-by: Codex --- test/sequential/test-pipe.js | 10 ++++------ 1 file changed, 4 insertions(+), 6 deletions(-) diff --git a/test/sequential/test-pipe.js b/test/sequential/test-pipe.js index 6f4a822ccf1b..612cf17a8698 100644 --- a/test/sequential/test-pipe.js +++ b/test/sequential/test-pipe.js @@ -25,8 +25,6 @@ const assert = require('assert'); const http = require('http'); const net = require('net'); -const webPort = common.PORT; -const tcpPort = webPort + 1; const bufferSize = 5 * 1024 * 1024; let listenCount = 0; @@ -45,7 +43,7 @@ const web = new http.Server(common.mustCall((req, res) => { web.close(); const socket = net.Stream(); - socket.connect(tcpPort); + socket.connect(tcp.address().port); socket.on('connect', common.mustCall()); @@ -60,7 +58,7 @@ const web = new http.Server(common.mustCall((req, res) => { req.connection.on('error', common.mustNotCall()); })); -web.listen(webPort, startClient); +web.listen(0, startClient); const tcp = net.Server(common.mustCall((s) => { @@ -83,14 +81,14 @@ const tcp = net.Server(common.mustCall((s) => { s.on('error', common.mustNotCall()); })); -tcp.listen(tcpPort, startClient); +tcp.listen(0, startClient); function startClient() { listenCount++; if (listenCount < 2) return; const req = http.request({ - port: common.PORT, + port: web.address().port, method: 'GET', path: '/', headers: { From ecc5e9250f830e5c84929e01ce4c94f62e5f12f1 Mon Sep 17 00:00:00 2001 From: Filip Skokan Date: Thu, 1 Oct 2026 16:31:20 +0200 Subject: [PATCH 10/10] test: isolate sequential worker port ranges Sequential workers share common.PORT even when different subsystems overlap. Assign 100-port ranges through NODE_COMMON_PORT, keeping the same range through retries and preserving the -j1 environment. Report the assigned base on failure and reject configurations that exceed the port range. Signed-off-by: Filip Skokan Assisted-by: Codex --- test/README.md | 9 +- test/tools/test_test_runner.py | 148 +++++++++++++++++++++++++++++++++ tools/test.py | 37 +++++++++ 3 files changed, 192 insertions(+), 2 deletions(-) diff --git a/test/README.md b/test/README.md index 42a1cbe3ae1e..3bd95329ea49 100644 --- a/test/README.md +++ b/test/README.md @@ -53,8 +53,13 @@ The test runner finishes parallel tests first, followed by serial suites, then up to the worker count selected with `-j`. The subsystem is the first filename component after `test-`, so `test-net-server-bind.js` and `test-net-connect-econnrefused.js` cannot overlap, while a `test-fs-*` test can run -alongside them. This scheduling does not isolate resources shared across -subsystems, such as fixed ports. +alongside them. Each worker gets a separate range of 100 ports through +`NODE_COMMON_PORT`, starting at 12346 or the configured `NODE_COMMON_PORT` base. +Tests and their child processes retain that range through retries. With `-j1`, +the existing port environment is preserved. Multi-worker runs require a base +that leaves room for every worker's range below port 65536. +Ports specified independently of `common.PORT` and other shared resources +still require isolation. [^1]: [Documentation](../test/common/README.md) diff --git a/test/tools/test_test_runner.py b/test/tools/test_test_runner.py index 9bfd3ba172eb..5ca0f9bdf903 100644 --- a/test/tools/test_test_runner.py +++ b/test/tools/test_test_runner.py @@ -1,6 +1,8 @@ import contextlib +import copy import io import os +import socket import sys import threading import unittest @@ -59,6 +61,27 @@ def HasRun(self, output): event.set() +class EnvironmentCase(SchedulerCase): + def __init__(self, name, action=None, suite='sequential', parallel=False): + super().__init__(name, action, suite, parallel) + self.environments = [] + + def Run(self): + self.calls += 1 + return runner.TestCase.Run(self) + + def GetRunConfiguration(self): + return {'command': ['node', '/'.join(self.path)], 'envs': {'EXAMPLE': 'kept'}} + + def GetLabel(self): + return '/'.join(self.path) + + def RunCommand(self, command, env): + self.environments.append(env.copy()) + output = self.action(self, env) if self.action else None + return runner.TestOutput(self, command, output or runner.CommandOutput(0, False, '', ''), False) + + class ExecuteTest(unittest.TestCase): def test_invalid_utf8_output_preserves_test_result(self): context = SimpleNamespace(verbose=False, suppress_dialogs=False, @@ -263,6 +286,131 @@ def test_assigns_unique_serial_ids_and_bounded_worker_ids(self): self.assertEqual(len(progress.completed), len(cases)) self.assertTrue(all(case.duration is not None for case in cases)) + def test_worker_port_ranges_allow_distinct_subsystem_listeners_to_overlap(self): + started = [threading.Event(), threading.Event()] + release = threading.Event() + # Select a base and check that its adjacent worker range is available. + for _ in range(20): + with socket.socket() as first, socket.socket() as second: + first.bind(('127.0.0.1', 0)) + base = first.getsockname()[1] + if base + 2 * runner.SEQUENTIAL_PORT_RANGE - 1 > 65535: + continue + try: + second.bind(('127.0.0.1', base + runner.SEQUENTIAL_PORT_RANGE)) + except OSError: + continue + break + else: + self.fail('could not find two available worker port ranges') + + def listener(index): + def action(case, env): + with socket.socket() as server: + port = int(env.get('NODE_COMMON_PORT', os.environ['NODE_COMMON_PORT'])) + server.bind(('127.0.0.1', port)) + server.listen() + started[index].set() + self.wait_for(started[1 - index]) + self.wait_for(release) + return action + + cases = [EnvironmentCase('test-net-port', listener(0)), + EnvironmentCase('test-http-port', listener(1))] + progress = SchedulerProgress(cases) + with mock.patch.dict(os.environ, {'NODE_COMMON_PORT': str(base)}): + with self.running(progress, release=[release]) as complete: + for event in started: + self.wait_for(event) + release.set() + self.assertTrue(complete()['allPassed']) + ports = [int(case.environments[0]['NODE_COMMON_PORT']) for case in cases] + self.assertEqual(sorted(ports), [base, base + runner.SEQUENTIAL_PORT_RANGE]) + self.assertTrue(all(env['EXAMPLE'] == 'kept' for case in cases for env in case.environments)) + + def test_port_ranges_preserve_retries_repeats_and_other_suites(self): + def retry(case, env): + if case.calls == 1: + return runner.CommandOutput(1, False, '', '') + + retried = EnvironmentCase('test-net-retry', retry) + retried.outcomes.add(runner.FLAKY) + repeated = copy.deepcopy(retried) + parallel = EnvironmentCase('test-parallel', suite='parallel', parallel=True) + serial = EnvironmentCase('test-serial', suite='pummel') + progress = SchedulerProgress([parallel, serial, retried, repeated], runner.KEEP_RETRYING) + with mock.patch.dict(os.environ, {'NODE_COMMON_PORT': '20000'}), \ + contextlib.redirect_stdout(io.StringIO()): + self.assertTrue(progress.Run(4)['allPassed']) + for case in [retried, repeated]: + self.assertEqual(case.calls, 2) + self.assertEqual([env['NODE_COMMON_PORT'] for env in case.environments], + [str(20000 + case.thread_id * runner.SEQUENTIAL_PORT_RANGE)] * 2) + self.assertTrue(all(env['TEST_PARALLEL'] == '0' for env in case.environments)) + for case in [parallel, serial]: + self.assertNotIn('NODE_COMMON_PORT', case.environments[0]) + + def test_single_worker_preserves_port_environment(self): + case = EnvironmentCase('test-net-port') + with mock.patch.dict(os.environ, {'NODE_COMMON_PORT': 'invalid'}), \ + mock.patch.object(runner.threading, 'Thread') as thread: + self.assertTrue(SchedulerProgress([case]).Run(1)['allPassed']) + thread.assert_not_called() + self.assertNotIn('NODE_COMMON_PORT', case.environments[0]) + + def test_port_range_defaults_and_valid_boundary(self): + for configured, expected in [('', runner.DEFAULT_COMMON_PORT), + ('0', runner.DEFAULT_COMMON_PORT), + ('65136', 65136)]: + with self.subTest(configured=configured), \ + mock.patch.dict(os.environ, {'NODE_COMMON_PORT': configured}): + case = EnvironmentCase('test-net-port') + self.assertTrue(SchedulerProgress([case]).Run(4)['allPassed']) + self.assertEqual(case.environments[0]['NODE_COMMON_PORT'], + str(expected + case.thread_id * runner.SEQUENTIAL_PORT_RANGE)) + + def test_invalid_port_ranges_fail_before_starting_any_tests(self): + for base in ['invalid', '-1', '65137']: + with self.subTest(base=base), mock.patch.dict(os.environ, {'NODE_COMMON_PORT': base}), \ + mock.patch.object(runner.threading, 'Thread') as thread: + cases = [EnvironmentCase('test-net-port'), + EnvironmentCase('test-parallel', suite='parallel', parallel=True)] + progress = SchedulerProgress(cases) + with self.assertRaises(runner.PortRangeError): + progress.Run(4) + thread.assert_not_called() + self.assertEqual([case.calls for case in cases], [0, 0]) + self.assertTrue(progress.shutdown_event.is_set()) + + with mock.patch.dict(os.environ, {'NODE_COMMON_PORT': 'invalid'}): + case = EnvironmentCase('test-serial', suite='pummel') + self.assertTrue(SchedulerProgress([case]).Run(4)['allPassed']) + + def test_failure_report_includes_assigned_port(self): + def fail(case, env): + return runner.CommandOutput(1, False, '', '') + + case = EnvironmentCase('test-net-failed', fail) + with mock.patch.dict(os.environ, {'NODE_COMMON_PORT': '20000'}): + progress = SchedulerProgress([case]) + self.assertFalse(progress.Run(2)['allPassed']) + self.assertIn('Environment: NODE_COMMON_PORT=%d' % case.common_port, + progress.GetFailureOutput(progress.failed[0])) + + failure = progress.failed[0] + for indicator in [runner.MonochromeProgressIndicator, runner.ColorProgressIndicator]: + with self.subTest(indicator=indicator.__name__), \ + contextlib.redirect_stdout(io.StringIO()) as report: + indicator([case], runner.RUN, 0).HasRun(failure) + self.assertIn('Environment: NODE_COMMON_PORT=%d' % case.common_port, report.getvalue()) + + with mock.patch.object(runner.logger, 'info') as log: + tap = runner.TapProgressIndicator([case], runner.RUN, 0) + tap.Starting() + tap.HasRun(failure) + self.assertIn(mock.call(' environment: {NODE_COMMON_PORT: %d}', case.common_port), + log.call_args_list) + def check_retries(self, measure_flakiness): retry_started = threading.Event() other_reported = threading.Event() diff --git a/tools/test.py b/tools/test.py index f8bfbb8f6cf2..a0ea4df3ff20 100755 --- a/tools/test.py +++ b/tools/test.py @@ -87,6 +87,15 @@ def get_module(name, path): VERBOSE = False +# Inspector tests use offsets through common.PORT + 92. +SEQUENTIAL_PORT_RANGE = 100 +DEFAULT_COMMON_PORT = 12346 + + +class PortRangeError(ValueError): + pass + + os.umask(0o022) os.environ.pop('NODE_OPTIONS', None) @@ -106,6 +115,7 @@ def __init__(self, cases, flaky_tests_mode, measure_flakiness): self.serial_queue = Queue(len(cases)) self.sequential_queue = [] self.running_subsystems = set() + self.sequential_port_base = None self.condition = threading.Condition() for case in cases: if case.parallel: @@ -133,6 +143,8 @@ def GetFailureOutput(self, failure): output += ["--- stdout ---"] output += [failure.output.stdout.strip()] output += ["Command: %s" % failure.test.GetFailureCommand(failure.command)] + if failure.test.common_port is not None: + output += ["Environment: NODE_COMMON_PORT=%d" % failure.test.common_port] if failure.HasCrashed(): output += ["--- %s ---" % PrintCrashed(failure.output.exit_code)] if failure.HasTimedOut(): @@ -157,6 +169,18 @@ def PrintFailureHeader(self, test): def Run(self, tasks) -> Dict: self.Starting() try: + if self.sequential_queue and tasks > 1: + try: + self.sequential_port_base = int(os.environ.get('NODE_COMMON_PORT') or + DEFAULT_COMMON_PORT) or DEFAULT_COMMON_PORT + except ValueError as error: + raise PortRangeError('NODE_COMMON_PORT must be an integer') from error + if (self.sequential_port_base < 1 or + self.sequential_port_base + tasks * SEQUENTIAL_PORT_RANGE - 1 > 65535): + raise PortRangeError( + 'NODE_COMMON_PORT=%d cannot provide %d worker ranges of %d ports; ' + 'lower NODE_COMMON_PORT or -j' % + (self.sequential_port_base, tasks, SEQUENTIAL_PORT_RANGE)) self.RunPhase(self.parallel_queue, tasks) self.RunSingle(self.serial_queue, 0) self.RunPhase(self.sequential_queue, tasks) @@ -271,6 +295,9 @@ def RunSingle(self, queue, thread_id): def RunCase(self, case, thread_id): with self.lock: case.thread_id = thread_id + case.common_port = (self.sequential_port_base + thread_id * SEQUENTIAL_PORT_RANGE + if case.path[0] == 'sequential' and + self.sequential_port_base is not None else None) case.serial_id = self.serial_id self.serial_id += 1 self.AboutToRun(case) @@ -494,6 +521,8 @@ def HasRun(self, output): duration = output.test.duration logger.info(' ---') logger.info(' duration_ms: %.5f' % (duration / timedelta(milliseconds=1))) + if output.UnexpectedOutput() and output.test.common_port is not None: + logger.info(' environment: {NODE_COMMON_PORT: %d}', output.test.common_port) if self.severity != 'ok' or self.traceback != '': if output.HasTimedOut(): self.traceback = 'timeout\n' + output.output.stdout + output.output.stderr @@ -562,6 +591,8 @@ def HasRun(self, output): if len(stderr): print(self.templates['stderr'] % stderr) print("Command: %s" % output.test.GetFailureCommand(output.command)) + if output.test.common_port is not None: + print("Environment: NODE_COMMON_PORT=%d" % output.test.common_port) if output.HasCrashed(): print("--- %s ---" % PrintCrashed(output.output.exit_code)) if output.HasTimedOut(): @@ -660,6 +691,7 @@ def __init__(self, context, path, arch, mode): self.max_virtual_memory = None self.serial_id = 0 self.thread_id = 0 + self.common_port = None def GetReportingName(self, command): prefix = abspath(join(dirname(__file__), '../test')) + os.sep @@ -708,6 +740,8 @@ def Run(self): "TEST_PARALLEL" : "%d" % self.parallel, "GITHUB_STEP_SUMMARY": "", }) + if self.common_port is not None: + envs['NODE_COMMON_PORT'] = str(self.common_port) result = self.RunCommand( command, envs @@ -1932,6 +1966,9 @@ def should_keep(case): except KeyboardInterrupt: print("Interrupted") return 1 + except PortRangeError as error: + PrintError(str(error)) + return 1 if options.time: # Write the times to stderr to make it easy to separate from the