Compare commits

..

23 Commits

Author SHA1 Message Date
Aleksandr Mezin
b591be4868 Release 3.2.1
Some checks failed
deploy / deploy (push) Failing after 1s
2023-03-12 14:42:09 +02:00
Aleksandr Mezin
b6c195a56c Merge pull request #885 from amezin/steal-hang-fix
Fix hang caused by `steal` command with empty test queue
2023-03-10 13:12:30 +02:00
Aleksandr Mezin
6abcdfc22e Fix hang caused by steal command with empty test queue
Fixes #884
2023-03-09 16:44:34 +02:00
pre-commit-ci[bot]
58fd7ccc05 [pre-commit.ci] pre-commit autoupdate (#881)
updates:
- [github.com/pre-commit/mirrors-mypy: v1.0.0 → v1.0.1](https://github.com/pre-commit/mirrors-mypy/compare/v1.0.0...v1.0.1)

Co-authored-by: pre-commit-ci[bot] <66853113+pre-commit-ci[bot]@users.noreply.github.com>
2023-02-27 08:53:38 -03:00
Bruno Oliveira
ba526fad5a Merge pull request #879 from pytest-dev/pre-commit-ci-update-config
[pre-commit.ci] pre-commit autoupdate
2023-02-14 07:29:00 -03:00
pre-commit-ci[bot]
efe674b265 [pre-commit.ci] pre-commit autoupdate
updates:
- [github.com/pre-commit/mirrors-mypy: v0.991 → v1.0.0](https://github.com/pre-commit/mirrors-mypy/compare/v0.991...v1.0.0)
2023-02-14 03:23:39 +00:00
Bruno Oliveira
5d692a7d63 Merge pull request #878 from akx/patch-1
docs: Remove unused statement in one-log-per-worker example
2023-02-13 11:01:58 -03:00
Aarni Koskela
5e795d88e7 docs: Remove unused statement in one-log-per-worker example 2023-02-13 15:38:59 +02:00
Aleksandr Mezin
2329d3454f Merge pull request #875 from pytest-dev/release-3.2.0
Release 3.2.0
2023-02-07 17:08:32 +02:00
Aleksandr Mezin
5c065198e9 Release 3.2.0
Some checks failed
deploy / deploy (push) Failing after 2s
Fixes #874
2023-02-07 16:46:47 +02:00
pre-commit-ci[bot]
c695763e92 [pre-commit.ci] pre-commit autoupdate (#869)
* [pre-commit.ci] pre-commit autoupdate

updates:
- [github.com/PyCQA/autoflake: v2.0.0 → v2.0.1](https://github.com/PyCQA/autoflake/compare/v2.0.0...v2.0.1)
- [github.com/psf/black: 22.12.0 → 23.1.0](https://github.com/psf/black/compare/22.12.0...23.1.0)
- [github.com/asottile/blacken-docs: v1.12.1 → 1.13.0](https://github.com/asottile/blacken-docs/compare/v1.12.1...1.13.0)

* Update .pre-commit-config.yaml

* [pre-commit.ci] auto fixes from pre-commit.com hooks

for more information, see https://pre-commit.ci

---------

Co-authored-by: pre-commit-ci[bot] <66853113+pre-commit-ci[bot]@users.noreply.github.com>
Co-authored-by: Bruno Oliveira <nicoddemus@gmail.com>
2023-02-07 07:58:55 -03:00
Bruno Oliveira
9373ddb3f8 Merge pull request #870 from hroncok/delenv
Tests: Unset PYTEST_XDIST_AUTO_NUM_WORKERS when the behavior without the envvar is asserted
2023-01-19 20:01:41 -03:00
Miro Hrončok
8fd1bfd182 Tests: Unset PYTEST_XDIST_AUTO_NUM_WORKERS when the behavior without the envvar is asserted 2023-01-17 20:10:12 +01:00
Thomas Kolar
017cc72b70 Document limitations for debugging as well as a workaround (#867)
Closes #863
2023-01-12 12:57:20 -03:00
Pat Thiel
cf19f76d86 Fix minor typos in the documentation (#866)
Some minor typo fixes.
2023-01-12 07:52:07 -03:00
Bruno Oliveira
691a0751bf Merge pull request #865 from amezin/fix-loadsched-tests
Fix some LoadScheduling tests
2023-01-11 20:34:27 -03:00
Aleksandr Mezin
e9860923bf Fix some LoadScheduling tests
The expected number of nodes didn't match throughout the test code.
2023-01-11 22:34:56 +02:00
Aleksandr Mezin
d1dfad3e92 Implement work-stealing scheduler (#862)
Closes #858
2023-01-11 08:38:57 -03:00
Aleksandr Mezin
9b0b5b1495 Add --maxschedchunk CLI option (#857)
Maximum number of tests scheduled in one step.

Setting it to 1 will force pytest to send tests to workers one by one -
might be useful for a small number of slow tests.

Larger numbers will allow the scheduler to submit consecutive chunks of tests
to workers - allows reusing fixtures.

Unlimited if not set.

Fixes #855
Fixes #255
2022-12-23 08:21:20 -03:00
pre-commit-ci[bot]
7faa69a04b [pre-commit.ci] pre-commit autoupdate (#856)
updates:
- [github.com/psf/black: 22.10.0 → 22.12.0](https://github.com/psf/black/compare/22.10.0...22.12.0)
- [github.com/asottile/pyupgrade: v3.3.0 → v3.3.1](https://github.com/asottile/pyupgrade/compare/v3.3.0...v3.3.1)

Co-authored-by: pre-commit-ci[bot] <66853113+pre-commit-ci[bot]@users.noreply.github.com>
2022-12-13 08:14:04 -03:00
Ronny Pfannschmidt
7a7be87453 Merge pull request #854 from pytest-dev/pre-commit-ci-update-config
[pre-commit.ci] pre-commit autoupdate
2022-12-06 10:47:31 +01:00
pre-commit-ci[bot]
d7877cef86 [pre-commit.ci] pre-commit autoupdate
updates:
- [github.com/asottile/pyupgrade: v3.2.2 → v3.3.0](https://github.com/asottile/pyupgrade/compare/v3.2.2...v3.3.0)
2022-12-06 00:37:27 +00:00
Bruno Oliveira
639e0868f7 Merge pull request #852 from nicoddemus/release-3.1.0
Release 3.1.0
2022-12-03 11:33:06 -03:00
17 changed files with 827 additions and 58 deletions

View File

@@ -1,19 +1,19 @@
repos: repos:
- repo: https://github.com/PyCQA/autoflake - repo: https://github.com/PyCQA/autoflake
rev: v2.0.0 rev: v2.0.1
hooks: hooks:
- id: autoflake - id: autoflake
args: ["--in-place", "--remove-unused-variables", "--remove-all-unused-imports"] args: ["--in-place", "--remove-unused-variables", "--remove-all-unused-imports"]
- repo: https://github.com/psf/black - repo: https://github.com/psf/black
rev: 22.10.0 rev: 23.1.0
hooks: hooks:
- id: black - id: black
args: [--safe, --quiet, --target-version, py35] args: [--safe, --quiet, --target-version, py35]
- repo: https://github.com/asottile/blacken-docs - repo: https://github.com/asottile/blacken-docs
rev: v1.12.1 rev: 1.13.0
hooks: hooks:
- id: blacken-docs - id: blacken-docs
additional_dependencies: [black==20.8b1] additional_dependencies: [black==23.1.0]
- repo: https://github.com/pre-commit/pre-commit-hooks - repo: https://github.com/pre-commit/pre-commit-hooks
rev: v4.4.0 rev: v4.4.0
hooks: hooks:
@@ -26,7 +26,7 @@ repos:
hooks: hooks:
- id: flake8 - id: flake8
- repo: https://github.com/asottile/pyupgrade - repo: https://github.com/asottile/pyupgrade
rev: v3.2.2 rev: v3.3.1
hooks: hooks:
- id: pyupgrade - id: pyupgrade
args: [--py3-plus] args: [--py3-plus]
@@ -39,7 +39,7 @@ repos:
language: python language: python
additional_dependencies: [pygments, restructuredtext_lint] additional_dependencies: [pygments, restructuredtext_lint]
- repo: https://github.com/pre-commit/mirrors-mypy - repo: https://github.com/pre-commit/mirrors-mypy
rev: v0.991 rev: v1.0.1
hooks: hooks:
- id: mypy - id: mypy
files: ^(src/|testing/) files: ^(src/|testing/)

View File

@@ -1,3 +1,36 @@
pytest-xdist 3.2.1 (2023-03-12)
===============================
Bug Fixes
---------
- `#884 <https://github.com/pytest-dev/pytest-xdist/issues/884>`_: Fixed hang in ``worksteal`` scheduler.
pytest-xdist 3.2.0 (2023-02-07)
===============================
Improved Documentation
----------------------
- `#863 <https://github.com/pytest-dev/pytest-xdist/issues/863>`_: Document limitations for debugging due to standard I/O of workers not being forwarded. Also, mention remote debugging as a possible workaround.
Features
--------
- `#855 <https://github.com/pytest-dev/pytest-xdist/issues/855>`_: Users can now configure ``load`` scheduling precision using ``--maxschedchunk`` command
line option.
- `#858 <https://github.com/pytest-dev/pytest-xdist/issues/858>`_: New ``worksteal`` scheduler, based on the idea of `work stealing <https://en.wikipedia.org/wiki/Work_stealing>`_. It's similar to ``load`` scheduler, but it should handle tests with significantly differing duration better, and, at the same time, it should provide similar or better reuse of fixtures.
Trivial Changes
---------------
- `#870 <https://github.com/pytest-dev/pytest-xdist/issues/870>`_: Make the tests pass even when ``$PYTEST_XDIST_AUTO_NUM_WORKERS`` is set.
pytest-xdist 3.1.0 (2022-12-01) pytest-xdist 3.1.0 (2022-12-01)
=============================== ===============================

View File

@@ -82,4 +82,13 @@ The test distribution algorithm is configured with the ``--dist`` command-line o
This will make sure ``test1`` and ``TestA::test2`` will run in the same worker. This will make sure ``test1`` and ``TestA::test2`` will run in the same worker.
Tests without the ``xdist_group`` mark are distributed normally as in the ``--dist=load`` mode. Tests without the ``xdist_group`` mark are distributed normally as in the ``--dist=load`` mode.
* ``--dist worksteal``: Initially, tests are distributed evenly among all
available workers. When a worker completes most of its assigned tests and
doesn't have enough tests to continue (currently, every worker needs at least
two tests in its queue), an attempt is made to reassign ("steal") a portion
of tests from some other worker's queue. The results should be similar to
the ``load`` method, but ``worksteal`` should handle tests with significantly
differing duration better, and, at the same time, it should provide similar
or better reuse of fixtures.
* ``--dist no``: The normal pytest execution mode, runs one test at a time (no distribution at all). * ``--dist no``: The normal pytest execution mode, runs one test at a time (no distribution at all).

View File

@@ -15,7 +15,7 @@ a test or fixture, you may use the ``worker_id`` fixture to do so:
@pytest.fixture() @pytest.fixture()
def user_account(worker_id): def user_account(worker_id):
""" use a different account in each xdist worker """ """use a different account in each xdist worker"""
return "account_%s" % worker_id return "account_%s" % worker_id
When ``xdist`` is disabled (running with ``-n0`` for example), then When ``xdist`` is disabled (running with ``-n0`` for example), then
@@ -80,7 +80,7 @@ wanted to create a separate database for each test run:
@pytest.fixture(scope="session", autouse=True) @pytest.fixture(scope="session", autouse=True)
def create_unique_database(testrun_uid): def create_unique_database(testrun_uid):
""" create a unique database for this particular test run """ """create a unique database for this particular test run"""
database_url = f"psql://myapp-{testrun_uid}" database_url = f"psql://myapp-{testrun_uid}"
with Semaphore(f"/{testrun_uid}-lock", flags=O_CREAT, initial_value=1): with Semaphore(f"/{testrun_uid}-lock", flags=O_CREAT, initial_value=1):
@@ -90,7 +90,7 @@ wanted to create a separate database for each test run:
@pytest.fixture() @pytest.fixture()
def db(testrun_uid): def db(testrun_uid):
""" retrieve unique database """ """retrieve unique database"""
database_url = f"psql://myapp-{testrun_uid}" database_url = f"psql://myapp-{testrun_uid}"
return database_get_instance(database_url) return database_get_instance(database_url)
@@ -221,7 +221,6 @@ Example:
def pytest_configure(config): def pytest_configure(config):
worker_id = os.environ.get("PYTEST_XDIST_WORKER") worker_id = os.environ.get("PYTEST_XDIST_WORKER")
if worker_id is not None: if worker_id is not None:
log_file = config.getini("worker_log_file")
logging.basicConfig( logging.basicConfig(
format=config.getini("log_file_format"), format=config.getini("log_file_format"),
filename=f"tests_{worker_id}.log", filename=f"tests_{worker_id}.log",

View File

@@ -6,7 +6,7 @@ pytest-xdist has some limitations that may be supported in pytest but can't be s
Order and amount of test must be consistent Order and amount of test must be consistent
------------------------------------------- -------------------------------------------
Is is not possible to have tests that differ in order or their amount across workers. It is not possible to have tests that differ in order or their amount across workers.
This is especially true with ``pytest.mark.parametrize``, when values are produced with sets or other unordered iterables/generators. This is especially true with ``pytest.mark.parametrize``, when values are produced with sets or other unordered iterables/generators.
@@ -59,6 +59,13 @@ Output (stdout and stderr) from workers
The ``-s``/``--capture=no`` option is meant to disable pytest capture, so users can then see stdout and stderr output in the terminal from tests and application code in real time. The ``-s``/``--capture=no`` option is meant to disable pytest capture, so users can then see stdout and stderr output in the terminal from tests and application code in real time.
However this option does not work with ``pytest-xdist`` because `execnet <https://github.com/pytest-dev/execnet>`__ the underlying library used for communication between master and workers, does not support transferring stdout/stderr from workers. However, this option does not work with ``pytest-xdist`` because `execnet <https://github.com/pytest-dev/execnet>`__ the underlying library used for communication between master and workers, does not support transferring stdout/stderr from workers.
Currenlty there are no plans ot support this in ``pytest-xdist``. Currently, there are no plans to support this in ``pytest-xdist``.
Debugging
~~~~~~~~~
This also means that debugging using PDB (or any other debugger that wants to use standard I/O) will not work. The ``--pdb`` option is disabled when distributing tests with ``pytest-xdist`` for this reason.
It is generally likely best to use ``pytest-xdist`` to find failing tests and then debug them without distribution; however, if you need to debug from within a worker process (for example, to address failures that only happen when running tests concurrently), remote debuggers (for example, `python-remote-pdb <https://github.com/ionelmc/python-remote-pdb>`__ or `python-web-pdb <https://github.com/romanvm/python-web-pdb>`__) have been reported to work for this purpose.

View File

@@ -8,6 +8,7 @@ from xdist.scheduler import (
LoadScopeScheduling, LoadScopeScheduling,
LoadFileScheduling, LoadFileScheduling,
LoadGroupScheduling, LoadGroupScheduling,
WorkStealingScheduling,
) )
@@ -100,6 +101,7 @@ class DSession:
"loadscope": LoadScopeScheduling, "loadscope": LoadScopeScheduling,
"loadfile": LoadFileScheduling, "loadfile": LoadFileScheduling,
"loadgroup": LoadGroupScheduling, "loadgroup": LoadGroupScheduling,
"worksteal": WorkStealingScheduling,
} }
return schedulers[dist](config, log) return schedulers[dist](config, log)
@@ -282,6 +284,17 @@ class DSession:
""" """
self.sched.mark_test_complete(node, item_index, duration) self.sched.mark_test_complete(node, item_index, duration)
def worker_unscheduled(self, node, indices):
"""
Emitted when a node fires the 'unscheduled' event, signalling that
some tests have been removed from the worker's queue and should be
sent to some worker again.
This should happen only in response to 'steal' command, so schedulers
not using 'steal' command don't have to implement it.
"""
self.sched.remove_pending_tests_from_node(node, indices)
def worker_collectreport(self, node, rep): def worker_collectreport(self, node, rep):
"""Emitted when a node calls the pytest_collectreport hook. """Emitted when a node calls the pytest_collectreport hook.

View File

@@ -35,7 +35,6 @@ def pytest_addoption(parser):
@pytest.hookimpl @pytest.hookimpl
def pytest_cmdline_main(config): def pytest_cmdline_main(config):
if config.getoption("looponfail"): if config.getoption("looponfail"):
usepdb = config.getoption("usepdb", False) # a core option usepdb = config.getoption("usepdb", False) # a core option
if usepdb: if usepdb:

View File

@@ -94,7 +94,15 @@ def pytest_addoption(parser):
"--dist", "--dist",
metavar="distmode", metavar="distmode",
action="store", action="store",
choices=["each", "load", "loadscope", "loadfile", "loadgroup", "no"], choices=[
"each",
"load",
"loadscope",
"loadfile",
"loadgroup",
"worksteal",
"no",
],
dest="dist", dest="dist",
default="no", default="no",
help=( help=(
@@ -107,6 +115,8 @@ def pytest_addoption(parser):
"loadfile: load balance by sending test grouped by file" "loadfile: load balance by sending test grouped by file"
" to any available environment.\n\n" " to any available environment.\n\n"
"loadgroup: like load, but sends tests marked with 'xdist_group' to the same worker.\n\n" "loadgroup: like load, but sends tests marked with 'xdist_group' to the same worker.\n\n"
"worksteal: split the test suite between available environments,"
" then rebalance when any worker runs out of tests.\n\n"
"(default) no: run tests inprocess, don't distribute." "(default) no: run tests inprocess, don't distribute."
), ),
) )
@@ -153,6 +163,19 @@ def pytest_addoption(parser):
"on every test run." "on every test run."
), ),
) )
group.addoption(
"--maxschedchunk",
action="store",
type=int,
help=(
"Maximum number of tests scheduled in one step for --dist=load. "
"Setting it to 1 will force pytest to send tests to workers one by "
"one - might be useful for a small number of slow tests. "
"Larger numbers will allow the scheduler to submit consecutive "
"chunks of tests to workers - allows reusing fixtures. "
"Unlimited if not set."
),
)
parser.addini( parser.addini(
"rsyncdirs", "rsyncdirs",

View File

@@ -6,6 +6,7 @@
needs not to be installed in remote environments. needs not to be installed in remote environments.
""" """
import contextlib
import sys import sys
import os import os
import time import time
@@ -56,14 +57,31 @@ def worker_title(title):
class WorkerInteractor: class WorkerInteractor:
SHUTDOWN_MARK = object()
QUEUE_REPLACED_MARK = object()
def __init__(self, config, channel): def __init__(self, config, channel):
self.config = config self.config = config
self.workerid = config.workerinput.get("workerid", "?") self.workerid = config.workerinput.get("workerid", "?")
self.testrunuid = config.workerinput["testrunuid"] self.testrunuid = config.workerinput["testrunuid"]
self.log = Producer(f"worker-{self.workerid}", enabled=config.option.debug) self.log = Producer(f"worker-{self.workerid}", enabled=config.option.debug)
self.channel = channel self.channel = channel
self.torun = self._make_queue()
self.nextitem_index = None
config.pluginmanager.register(self) config.pluginmanager.register(self)
def _make_queue(self):
return self.channel.gateway.execmodel.queue.Queue()
def _get_next_item_index(self):
"""Gets the next item from test queue. Handles the case when the queue
is replaced concurrently in another thread.
"""
result = self.torun.get()
while result is self.QUEUE_REPLACED_MARK:
result = self.torun.get()
return result
def sendevent(self, name, **kwargs): def sendevent(self, name, **kwargs):
self.log("sending", name, kwargs) self.log("sending", name, kwargs)
self.channel.send((name, kwargs)) self.channel.send((name, kwargs))
@@ -92,38 +110,63 @@ class WorkerInteractor:
def pytest_collection(self, session): def pytest_collection(self, session):
self.sendevent("collectionstart") self.sendevent("collectionstart")
def handle_command(self, command):
if command is self.SHUTDOWN_MARK:
self.torun.put(self.SHUTDOWN_MARK)
return
name, kwargs = command
self.log("received command", name, kwargs)
if name == "runtests":
for i in kwargs["indices"]:
self.torun.put(i)
elif name == "runtests_all":
for i in range(len(self.session.items)):
self.torun.put(i)
elif name == "shutdown":
self.torun.put(self.SHUTDOWN_MARK)
elif name == "steal":
self.steal(kwargs["indices"])
def steal(self, indices):
indices = set(indices)
stolen = []
old_queue, self.torun = self.torun, self._make_queue()
def old_queue_get_nowait_noraise():
with contextlib.suppress(self.channel.gateway.execmodel.queue.Empty):
return old_queue.get_nowait()
for i in iter(old_queue_get_nowait_noraise, None):
if i in indices:
stolen.append(i)
else:
self.torun.put(i)
self.sendevent("unscheduled", indices=stolen)
old_queue.put(self.QUEUE_REPLACED_MARK)
@pytest.hookimpl @pytest.hookimpl
def pytest_runtestloop(self, session): def pytest_runtestloop(self, session):
self.log("entering main loop") self.log("entering main loop")
torun = [] self.channel.setcallback(self.handle_command, endmarker=self.SHUTDOWN_MARK)
while 1: self.nextitem_index = self._get_next_item_index()
try: while self.nextitem_index is not self.SHUTDOWN_MARK:
name, kwargs = self.channel.receive() self.run_one_test()
except EOFError:
return True
self.log("received command", name, kwargs)
if name == "runtests":
torun.extend(kwargs["indices"])
elif name == "runtests_all":
torun.extend(range(len(session.items)))
self.log("items to run:", torun)
# only run if we have an item and a next item
while len(torun) >= 2:
self.run_one_test(torun)
if name == "shutdown":
if torun:
self.run_one_test(torun)
break
return True return True
def run_one_test(self, torun): def run_one_test(self):
self.item_index = self.nextitem_index
self.nextitem_index = self._get_next_item_index()
items = self.session.items items = self.session.items
self.item_index = torun.pop(0)
item = items[self.item_index] item = items[self.item_index]
if torun: if self.nextitem_index is self.SHUTDOWN_MARK:
nextitem = items[torun[0]]
else:
nextitem = None nextitem = None
else:
nextitem = items[self.nextitem_index]
worker_title("[pytest-xdist running] %s" % item.nodeid) worker_title("[pytest-xdist running] %s" % item.nodeid)

View File

@@ -3,3 +3,4 @@ from xdist.scheduler.load import LoadScheduling # noqa
from xdist.scheduler.loadfile import LoadFileScheduling # noqa from xdist.scheduler.loadfile import LoadFileScheduling # noqa
from xdist.scheduler.loadscope import LoadScopeScheduling # noqa from xdist.scheduler.loadscope import LoadScopeScheduling # noqa
from xdist.scheduler.loadgroup import LoadGroupScheduling # noqa from xdist.scheduler.loadgroup import LoadGroupScheduling # noqa
from xdist.scheduler.worksteal import WorkStealingScheduling # noqa

View File

@@ -64,6 +64,7 @@ class LoadScheduling:
else: else:
self.log = log.loadsched self.log = log.loadsched
self.config = config self.config = config
self.maxschedchunk = self.config.getoption("maxschedchunk")
@property @property
def nodes(self): def nodes(self):
@@ -185,7 +186,9 @@ class LoadScheduling:
# so let's rather wait with sending new items # so let's rather wait with sending new items
return return
num_send = items_per_node_max - len(node_pending) num_send = items_per_node_max - len(node_pending)
self._send_tests(node, num_send) # keep at least 2 tests pending even if --maxschedchunk=1
maxschedchunk = max(2 - len(node_pending), self.maxschedchunk)
self._send_tests(node, min(num_send, maxschedchunk))
else: else:
node.shutdown() node.shutdown()
@@ -245,6 +248,9 @@ class LoadScheduling:
if not self.collection: if not self.collection:
return return
if self.maxschedchunk is None:
self.maxschedchunk = len(self.collection)
# Send a batch of tests to run. If we don't have at least two # Send a batch of tests to run. If we don't have at least two
# tests per node, we have to send them all so that we can send # tests per node, we have to send them all so that we can send
# shutdown signals and get all nodes working. # shutdown signals and get all nodes working.
@@ -265,7 +271,8 @@ class LoadScheduling:
# how many items per node do we have about? # how many items per node do we have about?
items_per_node = len(self.collection) // len(self.node2pending) items_per_node = len(self.collection) // len(self.node2pending)
# take a fraction of tests for initial distribution # take a fraction of tests for initial distribution
node_chunksize = max(items_per_node // 4, 2) node_chunksize = min(items_per_node // 4, self.maxschedchunk)
node_chunksize = max(node_chunksize, 2)
# and initialize each node with a chunk of tests # and initialize each node with a chunk of tests
for node in self.nodes: for node in self.nodes:
self._send_tests(node, node_chunksize) self._send_tests(node, node_chunksize)

View File

@@ -213,13 +213,11 @@ class LoadScopeScheduling:
# A new node has been added later, perhaps an original one died. # A new node has been added later, perhaps an original one died.
if self.collection_is_completed: if self.collection_is_completed:
# Assert that .schedule() should have been called by now # Assert that .schedule() should have been called by now
assert self.collection assert self.collection
# Check that the new collection matches the official collection # Check that the new collection matches the official collection
if collection != self.collection: if collection != self.collection:
other_node = next(iter(self.registered_collections.keys())) other_node = next(iter(self.registered_collections.keys()))
msg = report_collection_diff( msg = report_collection_diff(

View File

@@ -0,0 +1,325 @@
from collections import namedtuple
from _pytest.runner import CollectReport
from xdist.remote import Producer
from xdist.workermanage import parse_spec_config
from xdist.report import report_collection_diff
NodePending = namedtuple("NodePending", ["node", "pending"])
# Every worker needs at least 2 tests in queue - the current and the next one.
MIN_PENDING = 2
class WorkStealingScheduling:
"""Implement work-stealing scheduling.
Initially, tests are distributed evenly among all nodes.
When some node completes most of its assigned tests (when only one pending
test remains), an attempt is made to reassign ("steal") some tests from
other nodes to this node.
Attributes:
:numnodes: The expected number of nodes taking part. The actual
number of nodes will vary during the scheduler's lifetime as
nodes are added by the DSession as they are brought up and
removed either because of a dead node or normal shutdown. This
number is primarily used to know when the initial collection is
completed.
:node2collection: Map of nodes and their test collection. All
collections should always be identical.
:node2pending: Map of nodes and the indices of their pending
tests. The indices are an index into ``.pending`` (which is
identical to their own collection stored in
``.node2collection``).
:collection: The one collection once it is validated to be
identical between all the nodes. It is initialised to None
until ``.schedule()`` is called.
:pending: List of indices of globally pending tests. These are
tests which have not yet been allocated to a chunk for a node
to process.
:log: A py.log.Producer instance.
:config: Config object, used for handling hooks.
:steal_requested_from_node: The node to which the current "steal" request
was sent. ``None`` if there is no request in progress. Only one request
can be in progress at any time, the scheduler doesn't send multiple
simultaneous requests.
"""
def __init__(self, config, log=None):
self.numnodes = len(parse_spec_config(config))
self.node2collection = {}
self.node2pending = {}
self.pending = []
self.collection = None
if log is None:
self.log = Producer("workstealsched")
else:
self.log = log.workstealsched
self.config = config
self.steal_requested_from_node = None
@property
def nodes(self):
"""A list of all nodes in the scheduler."""
return list(self.node2pending.keys())
@property
def collection_is_completed(self):
"""Boolean indication initial test collection is complete.
This is a boolean indicating all initial participating nodes
have finished collection. The required number of initial
nodes is defined by ``.numnodes``.
"""
return len(self.node2collection) >= self.numnodes
@property
def tests_finished(self):
"""Return True if all tests have been executed by the nodes."""
if not self.collection_is_completed:
return False
if self.pending:
return False
if self.steal_requested_from_node is not None:
return False
for pending in self.node2pending.values():
if len(pending) >= MIN_PENDING:
return False
return True
@property
def has_pending(self):
"""Return True if there are pending test items
This indicates that collection has finished and nodes are
still processing test items, so this can be thought of as
"the scheduler is active".
"""
if self.pending:
return True
for pending in self.node2pending.values():
if pending:
return True
return False
def add_node(self, node):
"""Add a new node to the scheduler.
From now on the node will be allocated chunks of tests to
execute.
Called by the ``DSession.worker_workerready`` hook when it
successfully bootstraps a new node.
"""
assert node not in self.node2pending
self.node2pending[node] = []
def add_node_collection(self, node, collection):
"""Add the collected test items from a node
The collection is stored in the ``.node2collection`` map.
Called by the ``DSession.worker_collectionfinish`` hook.
"""
assert node in self.node2pending
if self.collection_is_completed:
# A new node has been added later, perhaps an original one died.
# .schedule() should have
# been called by now
assert self.collection
if collection != self.collection:
other_node = next(iter(self.node2collection.keys()))
msg = report_collection_diff(
self.collection, collection, other_node.gateway.id, node.gateway.id
)
self.log(msg)
return
self.node2collection[node] = list(collection)
def mark_test_complete(self, node, item_index, duration=None):
"""Mark test item as completed by node
This is called by the ``DSession.worker_testreport`` hook.
"""
self.node2pending[node].remove(item_index)
self.check_schedule()
def mark_test_pending(self, item):
self.pending.insert(
0,
self.collection.index(item),
)
self.check_schedule()
def remove_pending_tests_from_node(self, node, indices):
"""Node returned some test indices back in response to 'steal' command.
This is called by ``DSession.worker_unscheduled``.
"""
assert node is self.steal_requested_from_node
self.steal_requested_from_node = None
indices_set = set(indices)
self.node2pending[node] = [
i for i in self.node2pending[node] if i not in indices_set
]
self.pending.extend(indices)
self.check_schedule()
def check_schedule(self):
"""Reschedule tests/perform load balancing."""
nodes_up = [
NodePending(node, pending)
for node, pending in self.node2pending.items()
if not node.shutting_down
]
def get_idle_nodes():
return [node for node, pending in nodes_up if len(pending) < MIN_PENDING]
idle_nodes = get_idle_nodes()
if not idle_nodes:
return
if self.pending:
# Distribute pending tests evenly among idle nodes
for i, node in enumerate(idle_nodes):
nodes_remaining = len(idle_nodes) - i
num_send = len(self.pending) // nodes_remaining
self._send_tests(node, num_send)
idle_nodes = get_idle_nodes()
# No need to steal anything if all nodes have enough work to continue
if not idle_nodes:
return
# Only one active stealing request is allowed
if self.steal_requested_from_node is not None:
return
# Find the node that has the longest test queue
steal_from = max(
nodes_up, key=lambda node_pending: len(node_pending.pending), default=None
)
if steal_from is None:
num_steal = 0
else:
# Steal half of the test queue - but keep that node running too.
# If the node has 2 or less tests queued, stealing will fail
# anyway.
max_steal = max(0, len(steal_from.pending) - MIN_PENDING)
num_steal = min(len(steal_from.pending) // 2, max_steal)
if num_steal == 0:
# Can't get more work - shutdown idle nodes. This will force them
# to run the last test now instead of waiting for more tests.
for node in idle_nodes:
node.shutdown()
return
steal_from.node.send_steal(steal_from.pending[-num_steal:])
self.steal_requested_from_node = steal_from.node
def remove_node(self, node):
"""Remove a node from the scheduler
This should be called either when the node crashed or at
shutdown time. In the former case any pending items assigned
to the node will be re-scheduled. Called by the
``DSession.worker_workerfinished`` and
``DSession.worker_errordown`` hooks.
Return the item which was being executing while the node
crashed or None if the node has no more pending items.
"""
pending = self.node2pending.pop(node)
# If node was removed without completing its assigned tests - it crashed
if pending:
crashitem = self.collection[pending.pop(0)]
else:
crashitem = None
self.pending.extend(pending)
# Dead node won't respond to "steal" request
if self.steal_requested_from_node is node:
self.steal_requested_from_node = None
self.check_schedule()
return crashitem
def schedule(self):
"""Initiate distribution of the test collection
Initiate scheduling of the items across the nodes. If this
gets called again later it behaves the same as calling
``.check_schedule()`` on all nodes so that newly added nodes
will start to be used.
This is called by the ``DSession.worker_collectionfinish`` hook
if ``.collection_is_completed`` is True.
"""
assert self.collection_is_completed
# Initial distribution already happened, reschedule on all nodes
if self.collection is not None:
self.check_schedule()
return
if not self._check_nodes_have_same_collection():
self.log("**Different tests collected, aborting run**")
return
# Collections are identical, create the index of pending items.
self.collection = list(self.node2collection.values())[0]
self.pending[:] = range(len(self.collection))
if not self.collection:
return
self.check_schedule()
def _send_tests(self, node, num):
tests_per_node = self.pending[:num]
if tests_per_node:
del self.pending[:num]
self.node2pending[node].extend(tests_per_node)
node.send_runtest_some(tests_per_node)
def _check_nodes_have_same_collection(self):
"""Return True if all nodes have collected the same items.
If collections differ, this method returns False while logging
the collection differences and posting collection errors to
pytest_collectreport hook.
"""
node_collection_items = list(self.node2collection.items())
first_node, col = node_collection_items[0]
same_collection = True
for node, collection in node_collection_items[1:]:
msg = report_collection_diff(
col, collection, first_node.gateway.id, node.gateway.id
)
if msg:
same_collection = False
self.log(msg)
if self.config is not None:
rep = CollectReport(
node.gateway.id, "failed", longrepr=msg, result=[]
)
self.config.hook.pytest_collectreport(report=rep)
return same_collection

View File

@@ -300,6 +300,9 @@ class WorkerController:
def send_runtest_all(self): def send_runtest_all(self):
self.sendcommand("runtests_all") self.sendcommand("runtests_all")
def send_steal(self, indices):
self.sendcommand("steal", indices=indices)
def shutdown(self): def shutdown(self):
if not self._down: if not self._down:
try: try:
@@ -359,6 +362,8 @@ class WorkerController:
self.notify_inproc(eventname, node=self, ids=kwargs["ids"]) self.notify_inproc(eventname, node=self, ids=kwargs["ids"])
elif eventname == "runtest_protocol_complete": elif eventname == "runtest_protocol_complete":
self.notify_inproc(eventname, node=self, **kwargs) self.notify_inproc(eventname, node=self, **kwargs)
elif eventname == "unscheduled":
self.notify_inproc(eventname, node=self, **kwargs)
elif eventname == "logwarning": elif eventname == "logwarning":
self.notify_inproc( self.notify_inproc(
eventname, eventname,

View File

@@ -1,6 +1,6 @@
from xdist.dsession import DSession, get_default_max_worker_restart from xdist.dsession import DSession, get_default_max_worker_restart
from xdist.report import report_collection_diff from xdist.report import report_collection_diff
from xdist.scheduler import EachScheduling, LoadScheduling from xdist.scheduler import EachScheduling, LoadScheduling, WorkStealingScheduling
from typing import Optional from typing import Optional
import pytest import pytest
@@ -17,6 +17,7 @@ class MockGateway:
class MockNode: class MockNode:
def __init__(self) -> None: def __init__(self) -> None:
self.sent = [] # type: ignore[var-annotated] self.sent = [] # type: ignore[var-annotated]
self.stolen = [] # type: ignore[var-annotated]
self.gateway = MockGateway() self.gateway = MockGateway()
self._shutdown = False self._shutdown = False
@@ -26,6 +27,9 @@ class MockNode:
def send_runtest_all(self) -> None: def send_runtest_all(self) -> None:
self.sent.append("ALL") self.sent.append("ALL")
def send_steal(self, indices) -> None:
self.stolen.extend(indices)
def shutdown(self) -> None: def shutdown(self) -> None:
self._shutdown = True self._shutdown = True
@@ -129,30 +133,79 @@ class TestLoadScheduling:
assert node1.sent == [0, 1, 4, 5] assert node1.sent == [0, 1, 4, 5]
assert not sched.pending assert not sched.pending
def test_schedule_fewer_tests_than_nodes(self, pytester: pytest.Pytester) -> None: def test_schedule_maxchunk_none(self, pytester: pytest.Pytester) -> None:
config = pytester.parseconfig("--tx=2*popen") config = pytester.parseconfig("--tx=2*popen")
sched = LoadScheduling(config) sched = LoadScheduling(config)
sched.add_node(MockNode()) sched.add_node(MockNode())
sched.add_node(MockNode()) sched.add_node(MockNode())
node1, node2 = sched.nodes
col = [f"test{i}" for i in range(16)]
sched.add_node_collection(node1, col)
sched.add_node_collection(node2, col)
sched.schedule()
assert node1.sent == [0, 1]
assert node2.sent == [2, 3]
assert sched.pending == list(range(4, 16))
assert sched.node2pending[node1] == node1.sent
assert sched.node2pending[node2] == node2.sent
sched.mark_test_complete(node1, 0)
assert node1.sent == [0, 1, 4, 5]
assert sched.pending == list(range(6, 16))
sched.mark_test_complete(node1, 1)
assert node1.sent == [0, 1, 4, 5]
assert sched.pending == list(range(6, 16))
for i in range(7, 16):
sched.mark_test_complete(node1, i - 3)
assert node1.sent == [0, 1] + list(range(4, i))
assert node2.sent == [2, 3]
assert sched.pending == list(range(i, 16))
def test_schedule_maxchunk_1(self, pytester: pytest.Pytester) -> None:
config = pytester.parseconfig("--tx=2*popen", "--maxschedchunk=1")
sched = LoadScheduling(config)
sched.add_node(MockNode())
sched.add_node(MockNode())
node1, node2 = sched.nodes
col = [f"test{i}" for i in range(16)]
sched.add_node_collection(node1, col)
sched.add_node_collection(node2, col)
sched.schedule()
assert node1.sent == [0, 1]
assert node2.sent == [2, 3]
assert sched.pending == list(range(4, 16))
assert sched.node2pending[node1] == node1.sent
assert sched.node2pending[node2] == node2.sent
for complete_index, first_pending in enumerate(range(5, 16)):
sched.mark_test_complete(node1, node1.sent[complete_index])
assert node1.sent == [0, 1] + list(range(4, first_pending))
assert node2.sent == [2, 3]
assert sched.pending == list(range(first_pending, 16))
def test_schedule_fewer_tests_than_nodes(self, pytester: pytest.Pytester) -> None:
config = pytester.parseconfig("--tx=3*popen")
sched = LoadScheduling(config)
sched.add_node(MockNode())
sched.add_node(MockNode())
sched.add_node(MockNode()) sched.add_node(MockNode())
node1, node2, node3 = sched.nodes node1, node2, node3 = sched.nodes
col = ["xyz"] * 2 col = ["xyz"] * 2
sched.add_node_collection(node1, col) sched.add_node_collection(node1, col)
sched.add_node_collection(node2, col) sched.add_node_collection(node2, col)
sched.add_node_collection(node3, col)
assert sched.collection_is_completed
sched.schedule() sched.schedule()
# assert not sched.tests_finished # assert not sched.tests_finished
sent1 = node1.sent assert node1.sent == [0]
sent2 = node2.sent assert node2.sent == [1]
sent3 = node3.sent assert node3.sent == []
assert sent1 == [0]
assert sent2 == [1]
assert sent3 == []
assert not sched.pending assert not sched.pending
def test_schedule_fewer_than_two_tests_per_node( def test_schedule_fewer_than_two_tests_per_node(
self, pytester: pytest.Pytester self, pytester: pytest.Pytester
) -> None: ) -> None:
config = pytester.parseconfig("--tx=2*popen") config = pytester.parseconfig("--tx=3*popen")
sched = LoadScheduling(config) sched = LoadScheduling(config)
sched.add_node(MockNode()) sched.add_node(MockNode())
sched.add_node(MockNode()) sched.add_node(MockNode())
@@ -161,14 +214,13 @@ class TestLoadScheduling:
col = ["xyz"] * 5 col = ["xyz"] * 5
sched.add_node_collection(node1, col) sched.add_node_collection(node1, col)
sched.add_node_collection(node2, col) sched.add_node_collection(node2, col)
sched.add_node_collection(node3, col)
assert sched.collection_is_completed
sched.schedule() sched.schedule()
# assert not sched.tests_finished # assert not sched.tests_finished
sent1 = node1.sent assert node1.sent == [0, 3]
sent2 = node2.sent assert node2.sent == [1, 4]
sent3 = node3.sent assert node3.sent == [2]
assert sent1 == [0, 3]
assert sent2 == [1, 4]
assert sent3 == [2]
assert not sched.pending assert not sched.pending
def test_add_remove_node(self, pytester: pytest.Pytester) -> None: def test_add_remove_node(self, pytester: pytest.Pytester) -> None:
@@ -217,6 +269,169 @@ class TestLoadScheduling:
assert "Different tests were collected between" in rep.longrepr assert "Different tests were collected between" in rep.longrepr
class TestWorkStealingScheduling:
def test_ideal_case(self, pytester: pytest.Pytester) -> None:
config = pytester.parseconfig("--tx=2*popen")
sched = WorkStealingScheduling(config)
sched.add_node(MockNode())
sched.add_node(MockNode())
node1, node2 = sched.nodes
collection = [f"test_workstealing.py::test_{i}" for i in range(16)]
assert not sched.collection_is_completed
sched.add_node_collection(node1, collection)
assert not sched.collection_is_completed
sched.add_node_collection(node2, collection)
assert sched.collection_is_completed
assert sched.node2collection[node1] == collection
assert sched.node2collection[node2] == collection
sched.schedule()
assert not sched.pending
assert not sched.tests_finished
assert node1.sent == list(range(0, 8))
assert node2.sent == list(range(8, 16))
for i in range(8):
sched.mark_test_complete(node1, node1.sent[i])
sched.mark_test_complete(node2, node2.sent[i])
assert sched.tests_finished
assert node1.stolen == []
assert node2.stolen == []
def test_stealing(self, pytester: pytest.Pytester) -> None:
config = pytester.parseconfig("--tx=2*popen")
sched = WorkStealingScheduling(config)
sched.add_node(MockNode())
sched.add_node(MockNode())
node1, node2 = sched.nodes
collection = [f"test_workstealing.py::test_{i}" for i in range(16)]
sched.add_node_collection(node1, collection)
sched.add_node_collection(node2, collection)
assert sched.collection_is_completed
sched.schedule()
assert node1.sent == list(range(0, 8))
assert node2.sent == list(range(8, 16))
for i in range(8):
sched.mark_test_complete(node1, node1.sent[i])
assert node2.stolen == list(range(12, 16))
sched.remove_pending_tests_from_node(node2, node2.stolen)
for i in range(4):
sched.mark_test_complete(node2, node2.sent[i])
assert node1.stolen == [14, 15]
sched.remove_pending_tests_from_node(node1, node1.stolen)
sched.mark_test_complete(node1, 12)
sched.mark_test_complete(node2, 14)
assert node2.stolen == list(range(12, 16))
assert node1.stolen == [14, 15]
assert sched.tests_finished
def test_steal_on_add_node(self, pytester: pytest.Pytester) -> None:
node = MockNode()
config = pytester.parseconfig("--tx=popen")
sched = WorkStealingScheduling(config)
sched.add_node(node)
collection = [f"test_workstealing.py::test_{i}" for i in range(5)]
sched.add_node_collection(node, collection)
assert sched.collection_is_completed
sched.schedule()
assert not sched.pending
sched.mark_test_complete(node, 0)
node2 = MockNode()
sched.add_node(node2)
sched.add_node_collection(node2, collection)
assert sched.collection_is_completed
sched.schedule()
assert node.stolen == [3, 4]
sched.remove_pending_tests_from_node(node, node.stolen)
sched.mark_test_complete(node, 1)
sched.mark_test_complete(node2, 3)
assert sched.tests_finished
assert node2.stolen == []
def test_schedule_fewer_tests_than_nodes(self, pytester: pytest.Pytester) -> None:
config = pytester.parseconfig("--tx=3*popen")
sched = WorkStealingScheduling(config)
sched.add_node(MockNode())
sched.add_node(MockNode())
sched.add_node(MockNode())
node1, node2, node3 = sched.nodes
col = ["xyz"] * 2
sched.add_node_collection(node1, col)
sched.add_node_collection(node2, col)
sched.add_node_collection(node3, col)
sched.schedule()
assert node1.sent == []
assert node1.stolen == []
assert node2.sent == [0]
assert node2.stolen == []
assert node3.sent == [1]
assert node3.stolen == []
assert not sched.pending
assert sched.tests_finished
def test_schedule_fewer_than_two_tests_per_node(
self, pytester: pytest.Pytester
) -> None:
config = pytester.parseconfig("--tx=3*popen")
sched = WorkStealingScheduling(config)
sched.add_node(MockNode())
sched.add_node(MockNode())
sched.add_node(MockNode())
node1, node2, node3 = sched.nodes
col = ["xyz"] * 5
sched.add_node_collection(node1, col)
sched.add_node_collection(node2, col)
sched.add_node_collection(node3, col)
sched.schedule()
assert node1.sent == [0]
assert node2.sent == [1, 2]
assert node3.sent == [3, 4]
assert not sched.pending
assert not sched.tests_finished
sched.mark_test_complete(node1, node1.sent[0])
sched.mark_test_complete(node2, node2.sent[0])
sched.mark_test_complete(node3, node3.sent[0])
sched.mark_test_complete(node3, node3.sent[1])
assert sched.tests_finished
assert node1.stolen == []
assert node2.stolen == []
assert node3.stolen == []
def test_add_remove_node(self, pytester: pytest.Pytester) -> None:
node = MockNode()
config = pytester.parseconfig("--tx=popen")
sched = WorkStealingScheduling(config)
sched.add_node(node)
collection = ["test_file.py::test_func"]
sched.add_node_collection(node, collection)
assert sched.collection_is_completed
sched.schedule()
assert not sched.pending
crashitem = sched.remove_node(node)
assert crashitem == collection[0]
def test_different_tests_collected(self, pytester: pytest.Pytester) -> None:
class CollectHook:
def __init__(self):
self.reports = []
def pytest_collectreport(self, report):
self.reports.append(report)
collect_hook = CollectHook()
config = pytester.parseconfig("--tx=2*popen")
config.pluginmanager.register(collect_hook, "collect_hook")
node1 = MockNode()
node2 = MockNode()
sched = WorkStealingScheduling(config)
sched.add_node(node1)
sched.add_node(node2)
sched.add_node_collection(node1, ["a.py::test_1"])
sched.add_node_collection(node2, ["a.py::test_2"])
sched.schedule()
assert len(collect_hook.reports) == 1
rep = collect_hook.reports[0]
assert "Different tests were collected between" in rep.longrepr
class TestDistReporter: class TestDistReporter:
@pytest.mark.xfail @pytest.mark.xfail
def test_rsync_printing(self, pytester: pytest.Pytester, linecomp) -> None: def test_rsync_printing(self, pytester: pytest.Pytester, linecomp) -> None:

View File

@@ -54,6 +54,8 @@ def test_auto_detect_cpus(
) -> None: ) -> None:
from xdist.plugin import pytest_cmdline_main as check_options from xdist.plugin import pytest_cmdline_main as check_options
monkeypatch.delenv("PYTEST_XDIST_AUTO_NUM_WORKERS", raising=False)
with suppress(ImportError): with suppress(ImportError):
import psutil import psutil
@@ -101,6 +103,7 @@ def test_auto_detect_cpus_psutil(
psutil = pytest.importorskip("psutil") psutil = pytest.importorskip("psutil")
monkeypatch.delenv("PYTEST_XDIST_AUTO_NUM_WORKERS", raising=False)
monkeypatch.setattr(psutil, "cpu_count", lambda logical=True: 84 if logical else 42) monkeypatch.setattr(psutil, "cpu_count", lambda logical=True: 84 if logical else 42)
config = pytester.parseconfigure("-nauto") config = pytester.parseconfigure("-nauto")
@@ -117,6 +120,8 @@ def test_auto_detect_cpus_os(
) -> None: ) -> None:
from xdist.plugin import pytest_cmdline_main as check_options from xdist.plugin import pytest_cmdline_main as check_options
monkeypatch.delenv("PYTEST_XDIST_AUTO_NUM_WORKERS", raising=False)
config = pytester.parseconfigure("-nauto") config = pytester.parseconfigure("-nauto")
check_options(config) check_options(config)
assert config.getoption("numprocesses") == 3 assert config.getoption("numprocesses") == 3
@@ -178,6 +183,8 @@ def test_hook_auto_num_workers_none(
# but we document it so let's test it. # but we document it so let's test it.
from xdist.plugin import pytest_cmdline_main as check_options from xdist.plugin import pytest_cmdline_main as check_options
monkeypatch.delenv("PYTEST_XDIST_AUTO_NUM_WORKERS", raising=False)
pytester.makeconftest( pytester.makeconftest(
""" """
def pytest_xdist_auto_num_workers(): def pytest_xdist_auto_num_workers():

View File

@@ -220,6 +220,91 @@ class TestWorkerInteractor:
ev = worker.popevent() ev = worker.popevent()
assert ev.name == "errordown" assert ev.name == "errordown"
def test_steal_work(self, worker: WorkerSetup, unserialize_report) -> None:
worker.pytester.makepyfile(
"""
import time
def test_func(): time.sleep(1)
def test_func2(): pass
def test_func3(): pass
def test_func4(): pass
"""
)
worker.setup()
ev = worker.popevent("collectionfinish")
ids = ev.kwargs["ids"]
assert len(ids) == 4
worker.sendcommand("runtests_all")
# wait for test_func setup
ev = worker.popevent("testreport")
rep = unserialize_report(ev.kwargs["data"])
assert rep.nodeid.endswith("::test_func")
assert rep.when == "setup"
worker.sendcommand("steal", indices=[1, 2])
ev = worker.popevent("unscheduled")
assert ev.kwargs["indices"] == [2]
reports = [
("test_func", "call"),
("test_func", "teardown"),
("test_func2", "setup"),
("test_func2", "call"),
("test_func2", "teardown"),
]
for func, when in reports:
ev = worker.popevent("testreport")
rep = unserialize_report(ev.kwargs["data"])
assert rep.nodeid.endswith(f"::{func}")
assert rep.when == when
worker.sendcommand("shutdown")
for when in ["setup", "call", "teardown"]:
ev = worker.popevent("testreport")
rep = unserialize_report(ev.kwargs["data"])
assert rep.nodeid.endswith("::test_func4")
assert rep.when == when
ev = worker.popevent("workerfinished")
assert "workeroutput" in ev.kwargs
def test_steal_empty_queue(self, worker: WorkerSetup, unserialize_report) -> None:
worker.pytester.makepyfile(
"""
def test_func(): pass
def test_func2(): pass
"""
)
worker.setup()
ev = worker.popevent("collectionfinish")
ids = ev.kwargs["ids"]
assert len(ids) == 2
worker.sendcommand("runtests_all")
for when in ["setup", "call", "teardown"]:
ev = worker.popevent("testreport")
rep = unserialize_report(ev.kwargs["data"])
assert rep.nodeid.endswith("::test_func")
assert rep.when == when
worker.sendcommand("steal", indices=[0, 1])
ev = worker.popevent("unscheduled")
assert ev.kwargs["indices"] == []
worker.sendcommand("shutdown")
for when in ["setup", "call", "teardown"]:
ev = worker.popevent("testreport")
rep = unserialize_report(ev.kwargs["data"])
assert rep.nodeid.endswith("::test_func2")
assert rep.when == when
ev = worker.popevent("workerfinished")
assert "workeroutput" in ev.kwargs
def test_remote_env_vars(pytester: pytest.Pytester) -> None: def test_remote_env_vars(pytester: pytest.Pytester) -> None:
pytester.makepyfile( pytester.makepyfile(