Compare commits

...

47 Commits

Author SHA1 Message Date
Ronny Pfannschmidt
771f248a04 refactor cpu number autodetection to avoid the regression 2015-08-19 21:34:26 +02:00
Ronny Pfannschmidt
c0c794b961 Added tag v1.13 for changeset 4e25f4c568be 2015-08-18 08:45:54 +02:00
Ronny Pfannschmidt
8630c27c75 flake8 fixes 2015-08-18 08:45:18 +02:00
Ronny Pfannschmidt
1fde875c91 changelog 2015-08-18 08:33:44 +02:00
Ronny Pfannschmidt
9f32c38398 add version file to hgignore 2015-08-18 08:30:00 +02:00
Ronny Pfannschmidt
eba5319fdb clean up setup.py and use setuptools_scm 2015-08-18 08:29:23 +02:00
Ronny Pfannschmidt
ec89a3c36b split up the plugin and extend tox test matrix 2015-08-08 11:57:08 +02:00
Ronny Pfannschmidt
6bab1dab75 simplify HostRsync constructor 2015-08-05 22:04:21 +02:00
Floris Bruynooghe
02cb571da0 Merged pull request #18 2015-08-04 23:30:50 +01:00
Floris Bruynooghe
f40e34c3ea Merged in nicoddemus/pytest-xdist/max-slave-restart-option (pull request #20)
Add --max-slave-restart option
2015-08-04 23:21:36 +01:00
Bruno Oliveira
4b1ddb9c81 Add --max-slave-restart option
Also changed wording used from "failed node" to "crashed slave", to conform
with other messages ("slave sw0 crashed")
2015-07-11 13:12:12 -03:00
Bruno Oliveira
9d4afbdfec Add test, CHANGELOG and docs for auto-cpu detection PR 2015-07-11 12:14:48 -03:00
Bruno Oliveira
ed9e5cd9ea merged auto-cpu branch 2015-07-11 12:01:00 -03:00
Bruno Oliveira
3330aac877 Add test for new pytest-2.8 behavior and add CHANGELOG entry 2015-07-11 11:45:27 -03:00
Bruno Oliveira
a131d34b2d testscollected is now a public Session attribute 2015-07-06 20:27:41 -03:00
Bruno Oliveira
05ab96feaf update collected tests from slaves into main pytest session object
Upstream changes in https://github.com/pytest-dev/pytest/pull/817 make
pytest-dist always return EXIT_NOTESTSCOLLECTED because the
master node doesn't collect any tests. This patch updates the pytest session
about the collected items on the slaves.
2015-07-04 15:40:55 -03:00
holger krekel
54a43db053 add support for releasing as universal wheel 2015-05-07 12:30:29 +02:00
holger krekel
1b0b406adb Added tag 1.12 for changeset 39ef85dbc893 2015-05-06 14:32:50 +02:00
holger krekel
2ea07f73a7 finalize 1.12 version, some more adaptation for pytest versions, streamlining tox.ini 2015-05-06 13:41:39 +02:00
holger krekel
faf2e0861f streamline tests so that they work wit pytest-2.8 2015-05-06 13:34:33 +02:00
holger krekel
eb53a5f8a0 Added tag 1.11 for changeset 220f6e46eb71 2015-04-16 08:48:06 +02:00
Anatoly Bubenkov
5da798f01f README.txt edited online with Bitbucket 2015-03-01 14:45:15 +00:00
holger krekel
94a7723ba8 fix link to pytest-xdist repository 2015-02-27 12:19:16 +01:00
Comrade DOS
9aee5f3d66 Add auto detection CPUs number. 2015-01-12 17:51:11 +06:00
holger krekel
87be3b7582 (added changelog) fix issue594: properly report errors when the test collection
is random.  Thanks Bruno Oliveira.
2014-09-24 13:43:57 +02:00
Bruno Oliveira
d84f1f08d8 fix issue 594: xdist is not executing tests parametrized with random values
Now xdist properly reports the collection errors instead of silently failing to execute
the test suite.
2014-09-23 22:09:44 -03:00
holger krekel
7a0cb35322 modernize tox.ini a bit and make another test dict-ordering independent 2014-09-18 19:40:48 +02:00
holger krekel
3a75b754d0 - add changelog entry for restart-crashed-nodes
- fix a test which depended on dict ordering
- minor test cleanups
2014-09-18 18:12:42 +02:00
Floris Bruynooghe
b418f0701c Functional tests for node restarting 2014-09-18 01:20:28 +01:00
Floris Bruynooghe
e3c22c65b8 Fix slavemanage tests for new NodeManager API 2014-09-18 00:48:34 +01:00
Floris Bruynooghe
c4e0327824 Give up on this one test 2014-09-17 23:43:07 +01:00
Floris Bruynooghe
9e28a56101 Fix rsync
Since this now gets called multiple times per gateway we need to
ensure an rsync happens for each combination of (spec, root) otherwise
only the first root for a gateway will be rsynced.
2014-09-17 23:42:33 +01:00
Floris Bruynooghe
76c297bdb0 Fixes for node-restarting in each scheduling
* node2collection should only be populated when the collection is
  added.  This way collection_is_completed can be kept track of
  correctly.

* The _removed2pending map needs to have the entire node as key
  because it is needed to look up the collection of the gateway which
  died in node2collection to check the collection of the replacement
  node is identical.

* init_distribute() needs to actually keep track of the started nodes
  otherwise it's (re)scheduling gets all confused and will re-schedule
  too often.
2014-09-16 23:18:49 +01:00
Floris Bruynooghe
fc335fb7dc Implement node-restarting for each scheduling 2014-09-16 16:01:51 +02:00
Floris Bruynooghe
f643991f7f Very rough cut of restarting a failed node 2014-09-15 01:20:45 +02:00
Floris Bruynooghe
969582187e Ignore pytest-cache's cache directory and pytestdebug.log 2014-09-13 11:06:43 +01:00
Anatoly Bubenkov
e275ce0f8e Merged in nicoddemus/pytest-xdist/log-collection-diff (pull request #9)
Log different tests collected by slaves instead of an error
2014-09-10 02:10:32 +02:00
Bruno Oliveira
f88275046f fixed tests removing line numbers from output
Pytest no longer show line numbers in the --verbose output
2014-09-09 20:28:16 -03:00
Bruno Oliveira
9d549a0c06 Log different tests collected by slaves instead of an error
This is a proposal to fix #556.
2014-08-02 20:59:45 -03:00
holger krekel
cc237b0d8e fix various flakes issues and add "flakes" to tox tests 2014-07-20 16:56:21 +02:00
holger krekel
67f80b7c2c fix pytest issue503: avoid random re-setup of broad scoped fixtures
(anything above function).
2014-07-20 16:41:03 +02:00
holger krekel
bfd00a58ff bump version, add py34 to tox.ini 2014-07-20 08:51:53 +02:00
holger krekel
ed12f57193 depend on latest py dev version to support "--boxed" properly 2014-07-05 17:34:56 +02:00
holger krekel
de377c1001 use config.getoption instead of deprecated config.getvalue
and adapt a few tests to also run against pytest-2.6.0.dev
2014-05-14 08:13:56 +02:00
holger krekel
eb639d19ec properly extend tests to cover xfailing capturing modes 2014-03-27 11:57:36 +01:00
holger krekel
4441002e95 fix pytest/xdist issue485 (also depends on py-1.4.21.dev1):
attach stdout/stderr on --boxed processes that die.
2014-03-26 18:33:09 +01:00
holger krekel
3f33a2d23e Added tag 1.10 for changeset 4406fc2a6427 2014-01-29 14:28:31 +01:00
22 changed files with 1288 additions and 593 deletions

View File

@@ -22,7 +22,10 @@ dist/
include/ include/
lib/ lib/
bin/ bin/
xdist/_version.py
pytest_xdist.egg-info pytest_xdist.egg-info
issue/ issue/
3rdparty/ 3rdparty/
pytestdebug.log
.tox .tox
.cache

View File

@@ -14,3 +14,7 @@ cd44a941c833c098e4899fe3d42a96703754d0d5 1.5
0d1c00018008433956aa7d93007bab6ea7de96e4 1.8 0d1c00018008433956aa7d93007bab6ea7de96e4 1.8
1d27987c267577899350a25ba5828d55d87083ad 1.8 1d27987c267577899350a25ba5828d55d87083ad 1.8
5c5cb6d59e12e566fbb0217aea718dc31578bee1 1.9 5c5cb6d59e12e566fbb0217aea718dc31578bee1 1.9
4406fc2a6427fadc021ed7e43e7aa5032b1ea91f 1.10
220f6e46eb71a6212ccbe6b67b9e6edcf8ee4fa5 1.11
39ef85dbc893cc63dede11601208098a667b58e9 1.12
4e25f4c568be2d7cb4d1739638a9e66bbf28f588 v1.13

View File

@@ -1,3 +1,59 @@
1.13.1
-------
- fix a regression -n 0 now disables xdist again
1.13
-------------------------
- extended the tox matrix with the supported py.test versions
- split up the plugin into 3 plugin's
to prepare the departure of boxed and looponfail.
looponfail will be a part of core
and forked boxed will be replaced
with a more reliable primitive based on xdist
- conforming with new pytest-2.8 behavior of returning non-zero when all
tests were skipped or deselected.
- new "--max-slave-restart" option that can be used to control maximum
number of times pytest-xdist can restart slaves due to crashes. Thanks to
Anatoly Bubenkov for the report and Bruno Oliveira for the PR.
- release as wheel
- "-n" option now can be set to "auto" for automatic detection of number
of cpus in the host system. Thanks Suloev Dmitry for the PR.
1.12
-------------------------
- fix issue594: properly report errors when the test collection
is random. Thanks Bruno Oliveira.
- some internal test suite adaptation (to become forward
compatible with the upcoming pytest-2.8)
1.11
-------------------------
- fix pytest/xdist issue485 (also depends on py-1.4.22):
attach stdout/stderr on --boxed processes that die.
- fix pytest/xdist issue503: make sure that a node has usually
two items to execute to avoid scoped fixtures to be torn down
pre-maturely (fixture teardown/setup is "nextitem" sensitive).
Thanks to Andreas Pelme for bug analysis and failing test.
- restart crashed nodes by internally refactoring setup handling
of nodes. Also includes better code documentation.
Many thanks to Floris Bruynooghe for the complete PR.
1.10 1.10
------------------------- -------------------------
@@ -120,4 +176,3 @@
- cleaned up termination handling - cleaned up termination handling
- make -x cause hard killing of test nodes to decrease wait time - make -x cause hard killing of test nodes to decrease wait time
until the traceback shows up on first failure until the traceback shows up on first failure

View File

@@ -1,5 +1,10 @@
.. image:: https://drone.io/bitbucket.org/pytest-dev/pytest-xdist/status.png
:target: https://drone.io/bitbucket.org/pytest-dev/pytest-xdist/latest
.. image:: https://pypip.in/v/pytest-xdist/badge.png
:target: https://pypi.python.org/pypi/pytest-xdist
xdist: pytest distributed testing plugin xdist: pytest distributed testing plugin
=============================================================== ============================
The `pytest-xdist`_ plugin extends py.test with some unique The `pytest-xdist`_ plugin extends py.test with some unique
test execution modes: test execution modes:
@@ -54,7 +59,13 @@ To send tests to multiple CPUs, type::
py.test -n NUM py.test -n NUM
Especially for longer running tests or tests requiring Especially for longer running tests or tests requiring
a lot of IO this can lead to considerable speed ups. a lot of IO this can lead to considerable speed ups. This option can
also be set to ``auto`` for automatic detection of the number of CPUs.
If a test crashes the interpreter, pytest-xdist will automatically restart
that slave and report the failure as usual. You can use the
``--max-slave-restart`` option to limit the number of slaves that can
be restarted, or disable restarting altogether using ``--max-slave-restart=0``.
Running tests in a Python subprocess Running tests in a Python subprocess
@@ -201,12 +212,10 @@ These directory specifications are relative to the directory
where the configuration file was found. where the configuration file was found.
.. _`pytest-xdist`: http://pypi.python.org/pypi/pytest-xdist .. _`pytest-xdist`: http://pypi.python.org/pypi/pytest-xdist
.. _`pytest-xdist repository`: http://bitbucket.org/hpk42/pytest-xdist .. _`pytest-xdist repository`: http://bitbucket.org/pytest-dev/pytest-xdist
.. _`pytest`: http://pytest.org .. _`pytest`: http://pytest.org
Issue and Bug Tracker Issue and Bug Tracker
------------------------ ------------------------
Please use the pytest issue tracker for bugs in this plugin, see https://bitbucket.org/hpk42/pytest/issues . Please use the pytest issue tracker for bugs in this plugin, see https://bitbucket.org/hpk42/pytest/issues .

2
setup.cfg Normal file
View File

@@ -0,0 +1,2 @@
[bdist_wheel]
universal = 1

View File

@@ -2,7 +2,7 @@ from setuptools import setup
setup( setup(
name="pytest-xdist", name="pytest-xdist",
version='1.10', use_scm_version={'write_to': 'xdist/_version.py'},
description='py.test xdist plugin for distributed testing and loop-on-failing modes', description='py.test xdist plugin for distributed testing and loop-on-failing modes',
long_description=open('README.txt').read(), long_description=open('README.txt').read(),
license='MIT', license='MIT',
@@ -11,20 +11,27 @@ setup(
url='http://bitbucket.org/hpk42/pytest-xdist', url='http://bitbucket.org/hpk42/pytest-xdist',
platforms=['linux', 'osx', 'win32'], platforms=['linux', 'osx', 'win32'],
packages = ['xdist'], packages = ['xdist'],
entry_points = {'pytest11': ['xdist = xdist.plugin'],}, entry_points = {
'pytest11': [
'xdist = xdist.plugin',
'xdist.looponfail = xdist.looponfail',
'xdist.boxed = xdist.boxed',
],
},
zip_safe=False, zip_safe=False,
install_requires = ['execnet>=1.1', 'pytest>=2.4.2'], install_requires=['execnet>=1.1', 'pytest>=2.4.2', 'py>=1.4.22'],
setup_requires=['setuptools_scm'],
classifiers=[ classifiers=[
'Development Status :: 5 - Production/Stable', 'Development Status :: 5 - Production/Stable',
'Intended Audience :: Developers', 'Intended Audience :: Developers',
'License :: OSI Approved :: MIT License', 'License :: OSI Approved :: MIT License',
'Operating System :: POSIX', 'Operating System :: POSIX',
'Operating System :: Microsoft :: Windows', 'Operating System :: Microsoft :: Windows',
'Operating System :: MacOS :: MacOS X', 'Operating System :: MacOS :: MacOS X',
'Topic :: Software Development :: Testing', 'Topic :: Software Development :: Testing',
'Topic :: Software Development :: Quality Assurance', 'Topic :: Software Development :: Quality Assurance',
'Topic :: Utilities', 'Topic :: Utilities',
'Programming Language :: Python', 'Programming Language :: Python',
'Programming Language :: Python :: 3', 'Programming Language :: Python :: 3',
], ],
) )

View File

@@ -1,5 +1,6 @@
import py import py
import sys import pytest
class TestDistribution: class TestDistribution:
def test_n1_pass(self, testdir): def test_n1_pass(self, testdir):
@@ -112,7 +113,7 @@ class TestDistribution:
import py import py
assert tmpdir.relto(py.path.local(%r)), tmpdir assert tmpdir.relto(py.path.local(%r)), tmpdir
""" % str(testdir.tmpdir)) """ % str(testdir.tmpdir))
result = testdir.runpytest(p1, "-n1") result = testdir.runpytest_subprocess(p1, "-n1")
assert result.ret == 0 assert result.ret == 0
result.stdout.fnmatch_lines([ result.stdout.fnmatch_lines([
"*1 passed*", "*1 passed*",
@@ -193,7 +194,7 @@ class TestDistribution:
assert dest.join(subdir.basename).check(dir=1) assert dest.join(subdir.basename).check(dir=1)
def test_data_exchange(self, testdir): def test_data_exchange(self, testdir):
c1 = testdir.makeconftest(""" testdir.makeconftest("""
# This hook only called on master. # This hook only called on master.
def pytest_configure_node(node): def pytest_configure_node(node):
node.slaveinput['a'] = 42 node.slaveinput['a'] = 42
@@ -242,7 +243,7 @@ class TestDistribution:
print ("s2call-finished") print ("s2call-finished")
""") """)
args = ["-n1", "--debug"] args = ["-n1", "--debug"]
result = testdir.runpytest(*args) result = testdir.runpytest_subprocess(*args)
s = result.stdout.str() s = result.stdout.str()
assert result.ret == 2 assert result.ret == 2
assert 's2call' in s assert 's2call' in s
@@ -250,14 +251,13 @@ class TestDistribution:
def test_keyboard_interrupt_dist(self, testdir): def test_keyboard_interrupt_dist(self, testdir):
# xxx could be refined to check for return code # xxx could be refined to check for return code
p = testdir.makepyfile(""" testdir.makepyfile("""
def test_sleep(): def test_sleep():
import time import time
time.sleep(10) time.sleep(10)
""") """)
child = testdir.spawn_pytest("-n1") child = testdir.spawn_pytest("-n1 -v")
py.std.time.sleep(0.1) child.expect(".*test_sleep.*")
child.expect(".*test session starts.*")
child.kill(2) # keyboard interrupt child.kill(2) # keyboard interrupt
child.expect(".*KeyboardInterrupt.*") child.expect(".*KeyboardInterrupt.*")
#child.expect(".*seconds.*") #child.expect(".*seconds.*")
@@ -270,7 +270,7 @@ class TestDistEach:
def test_hello(): def test_hello():
pass pass
""") """)
result = testdir.runpytest("--debug", "--dist=each", "--tx=2*popen") result = testdir.runpytest_subprocess("--debug", "--dist=each", "--tx=2*popen")
assert not result.ret assert not result.ret
result.stdout.fnmatch_lines(["*2 pass*"]) result.stdout.fnmatch_lines(["*2 pass*"])
@@ -300,7 +300,7 @@ class TestDistEach:
class TestTerminalReporting: class TestTerminalReporting:
def test_pass_skip_fail(self, testdir): def test_pass_skip_fail(self, testdir):
p = testdir.makepyfile(""" testdir.makepyfile("""
import py import py
def test_ok(): def test_ok():
pass pass
@@ -311,9 +311,9 @@ class TestTerminalReporting:
""") """)
result = testdir.runpytest("-n1", "-v") result = testdir.runpytest("-n1", "-v")
result.stdout.fnmatch_lines_random([ result.stdout.fnmatch_lines_random([
"*PASS*test_pass_skip_fail.py:2: *test_ok*", "*PASS*test_pass_skip_fail.py*test_ok*",
"*SKIP*test_pass_skip_fail.py:4: *test_skip*", "*SKIP*test_pass_skip_fail.py*test_skip*",
"*FAIL*test_pass_skip_fail.py:6: *test_func*", "*FAIL*test_pass_skip_fail.py*test_func*",
]) ])
result.stdout.fnmatch_lines([ result.stdout.fnmatch_lines([
"*def test_func():", "*def test_func():",
@@ -322,13 +322,13 @@ class TestTerminalReporting:
]) ])
def test_fail_platinfo(self, testdir): def test_fail_platinfo(self, testdir):
p = testdir.makepyfile(""" testdir.makepyfile("""
def test_func(): def test_func():
assert 0 assert 0
""") """)
result = testdir.runpytest("-n1", "-v") result = testdir.runpytest("-n1", "-v")
result.stdout.fnmatch_lines([ result.stdout.fnmatch_lines([
"*FAIL*test_fail_platinfo.py:1: *test_func*", "*FAIL*test_fail_platinfo.py*test_func*",
"*0*Python*", "*0*Python*",
"*def test_func():", "*def test_func():",
"> assert 0", "> assert 0",
@@ -351,7 +351,7 @@ def test_teardownfails_one_function(testdir):
@py.test.mark.xfail @py.test.mark.xfail
def test_terminate_on_hangingnode(testdir): def test_terminate_on_hangingnode(testdir):
p = testdir.makeconftest(""" p = testdir.makeconftest("""
def pytest_sessionfinishes(session): def pytest_sessionfinish(session):
if session.nodeid == "my": # running on slave if session.nodeid == "my": # running on slave
import time import time
time.sleep(3) time.sleep(3)
@@ -363,7 +363,20 @@ def test_terminate_on_hangingnode(testdir):
]) ])
def test_auto_detect_cpus(testdir, monkeypatch):
import multiprocessing
monkeypatch.setattr(multiprocessing, 'cpu_count', lambda: 3)
testdir.makeconftest("""
def pytest_unconfigure(config):
with open('cpus', 'w') as f:
f.write('cpus = %s' % config.option.numprocesses)
""")
testdir.inline_run('-n=auto')
cpus_file = testdir.tmpdir.join('cpus')
assert cpus_file.read() == 'cpus = 3'
@pytest.mark.xfail(reason="works if run outside test suite", run=False)
def test_session_hooks(testdir): def test_session_hooks(testdir):
testdir.makeconftest(""" testdir.makeconftest("""
import sys import sys
@@ -397,6 +410,31 @@ def test_session_hooks(testdir):
assert testdir.tmpdir.join("slave").check() assert testdir.tmpdir.join("slave").check()
assert testdir.tmpdir.join("master").check() assert testdir.tmpdir.join("master").check()
def test_session_testscollected(testdir):
"""
Make sure master node is updating the session object with the number
of tests collected from the slaves.
"""
testdir.makepyfile(test_foo="""
import pytest
@pytest.mark.parametrize('i', range(3))
def test_ok(i):
pass
""")
testdir.makeconftest("""
def pytest_sessionfinish(session):
collected = getattr(session, 'testscollected', None)
with open('testscollected', 'w') as f:
f.write('collected = %s' % collected)
""")
result = testdir.inline_run("-n1")
result.assertoutcome(passed=3)
collected_file = testdir.tmpdir.join('testscollected')
assert collected_file.isfile()
assert collected_file.read() == 'collected = 3'
def test_funcarg_teardown_failure(testdir): def test_funcarg_teardown_failure(testdir):
p = testdir.makepyfile(""" p = testdir.makepyfile("""
def pytest_funcarg__myarg(request): def pytest_funcarg__myarg(request):
@@ -407,7 +445,7 @@ def test_funcarg_teardown_failure(testdir):
def test_hello(myarg): def test_hello(myarg):
pass pass
""") """)
result = testdir.runpytest("--debug", p) # , "-n1") result = testdir.runpytest_subprocess("--debug", p) # , "-n1")
result.stdout.fnmatch_lines([ result.stdout.fnmatch_lines([
"*ValueError*42*", "*ValueError*42*",
"*1 passed*1 error*", "*1 passed*1 error*",
@@ -454,8 +492,142 @@ def test_issue34_pluginloading_in_subprocess(testdir):
def test_hello(): def test_hello():
assert pytest.sample_variable == "testing" assert pytest.sample_variable == "testing"
""") """)
result = testdir.runpytest("-n1", "-p", "plugin123") result = testdir.runpytest_subprocess("-n1", "-p", "plugin123")
assert result.ret == 0 assert result.ret == 0
result.stdout.fnmatch_lines([ result.stdout.fnmatch_lines([
"*1 passed*", "*1 passed*",
]) ])
def test_fixture_scope_caching_issue503(testdir):
p1 = testdir.makepyfile("""
import pytest
@pytest.fixture(scope='session')
def fix():
assert fix.counter == 0, 'session fixture was invoked multiple times'
fix.counter += 1
fix.counter = 0
def test_a(fix):
pass
def test_b(fix):
pass
""")
result = testdir.runpytest(p1, '-v', '-n1')
assert result.ret == 0
result.stdout.fnmatch_lines([
"*2 passed*",
])
def test_issue_594_random_parametrize(testdir):
"""
Make sure that tests that are randomly parametrized display an appropriate
error message, instead of silently skipping the entire test run.
"""
p1 = testdir.makepyfile("""
import pytest
import random
xs = list(range(10))
random.shuffle(xs)
@pytest.mark.parametrize('x', xs)
def test_foo(x):
assert 1
""")
result = testdir.runpytest(p1, '-v', '-n4')
assert result.ret == 1
result.stdout.fnmatch_lines([
"Different tests were collected between gw* and gw*",
])
class TestNodeFailure:
def test_load_single(self, testdir):
f = testdir.makepyfile("""
import os
def test_a(): os._exit(1)
def test_b(): pass
""")
res = testdir.runpytest(f, '-n1')
res.stdout.fnmatch_lines([
"*Replacing crashed slave*",
"*Slave*crashed while running*",
"*1 failed*1 passed*",
])
def test_load_multiple(self, testdir):
f = testdir.makepyfile("""
import os
def test_a(): pass
def test_b(): os._exit(1)
def test_c(): pass
def test_d(): pass
""")
res = testdir.runpytest(f, '-n2')
res.stdout.fnmatch_lines([
"*Replacing crashed slave*",
"*Slave*crashed while running*",
"*1 failed*3 passed*",
])
def test_each_single(self, testdir):
f = testdir.makepyfile("""
import os
def test_a(): os._exit(1)
def test_b(): pass
""")
res = testdir.runpytest(f, '--dist=each', '--tx=popen')
res.stdout.fnmatch_lines([
"*Replacing crashed slave*",
"*Slave*crashed while running*",
"*1 failed*1 passed*",
])
def test_each_multiple(self, testdir):
f = testdir.makepyfile("""
import os
def test_a(): os._exit(1)
def test_b(): pass
""")
res = testdir.runpytest(f, '--dist=each', '--tx=2*popen')
res.stdout.fnmatch_lines([
"*Replacing crashed slave*",
"*Slave*crashed while running*",
"*2 failed*2 passed*",
])
def test_max_slave_restart(self, testdir):
f = testdir.makepyfile("""
import os
def test_a(): pass
def test_b(): os._exit(1)
def test_c(): os._exit(1)
def test_d(): pass
""")
res = testdir.runpytest(f, '-n4', '--max-slave-restart=1')
res.stdout.fnmatch_lines([
"*Replacing crashed slave*",
"*Maximum crashed slaves reached: 1*",
"*Slave*crashed while running*",
"*Slave*crashed while running*",
"*2 failed*2 passed*",
])
def test_disable_restart(self, testdir):
f = testdir.makepyfile("""
import os
def test_a(): pass
def test_b(): os._exit(1)
def test_c(): pass
""")
res = testdir.runpytest(f, '-n4', '--max-slave-restart=0')
res.stdout.fnmatch_lines([
"*Slave restarting disabled*",
"*Slave*crashed while running*",
"*1 failed*2 passed*",
])

View File

@@ -1,20 +1,43 @@
import py import py
import pytest
import execnet import execnet
@pytest.fixture(scope="session", autouse=True)
def _ensure_imports():
# we import some modules because pytest-2.8's testdir fixture
# will unload all modules after each test and this cause
# (unknown) problems with execnet.Group()
execnet.Group
execnet.makegateway
pytest_plugins = "pytester" pytest_plugins = "pytester"
#rsyncdirs = ['.', '../xdist', py.path.local(execnet.__file__).dirpath()] #rsyncdirs = ['.', '../xdist', py.path.local(execnet.__file__).dirpath()]
@pytest.fixture(autouse=True)
def _divert_atexit(request, monkeypatch):
import atexit
l = []
def finish():
while l:
l.pop()()
monkeypatch.setattr(atexit, "register", l.append)
request.addfinalizer(finish)
def pytest_addoption(parser): def pytest_addoption(parser):
parser.addoption('--gx', parser.addoption('--gx',
action="append", dest="gspecs", default=None, action="append", dest="gspecs",
help=("add a global test environment, XSpec-syntax. ")) help=("add a global test environment, XSpec-syntax. "))
def pytest_funcarg__specssh(request): def pytest_funcarg__specssh(request):
return getspecssh(request.config) return getspecssh(request.config)
def getgspecs(config):
return [execnet.XSpec(spec) @pytest.fixture
for spec in config.getvalueorskip("gspecs")] def testdir(testdir):
# pytest before 2.8 did not have a runpytest_subprocess
if not hasattr(testdir, "runpytest_subprocess"):
testdir.runpytest_subprocess = testdir.runpytest
return testdir
# configuration information for tests # configuration information for tests
def getgspecs(config): def getgspecs(config):

View File

@@ -1,6 +1,10 @@
import py import pytest
import os
@py.test.mark.skipif("not hasattr(os, 'fork')") needsfork = pytest.mark.skipif(not hasattr(os, "fork"),
reason="os.fork required")
@needsfork
def test_functional_boxed(testdir): def test_functional_boxed(testdir):
p1 = testdir.makepyfile(""" p1 = testdir.makepyfile("""
import os import os
@@ -13,12 +17,36 @@ def test_functional_boxed(testdir):
"*1 failed*" "*1 failed*"
]) ])
@needsfork
@pytest.mark.parametrize("capmode", [
"no",
pytest.mark.xfail("sys", reason="capture cleanup needed"),
pytest.mark.xfail("fd", reason="capture cleanup needed")])
def test_functional_boxed_capturing(testdir, capmode):
p1 = testdir.makepyfile("""
import os
import sys
def test_function():
sys.stdout.write("hello\\n")
sys.stderr.write("world\\n")
os.kill(os.getpid(), 15)
""")
result = testdir.runpytest(p1, "--boxed", "--capture=%s" % capmode)
result.stdout.fnmatch_lines("""
*CRASHED*
*stdout*
hello
*stderr*
world
*1 failed*
""")
class TestOptionEffects: class TestOptionEffects:
def test_boxed_option_default(self, testdir): def test_boxed_option_default(self, testdir):
tmpdir = testdir.tmpdir.ensure("subdir", dir=1) tmpdir = testdir.tmpdir.ensure("subdir", dir=1)
config = testdir.parseconfig() config = testdir.parseconfig()
assert not config.option.boxed assert not config.option.boxed
py.test.importorskip("execnet") pytest.importorskip("execnet")
config = testdir.parseconfig('-d', tmpdir) config = testdir.parseconfig('-d', tmpdir)
assert not config.option.boxed assert not config.option.boxed

View File

@@ -4,7 +4,6 @@ from xdist.dsession import (
EachScheduling, EachScheduling,
report_collection_diff, report_collection_diff,
) )
from _pytest import main as outcome
import py import py
import pytest import pytest
import execnet import execnet
@@ -84,11 +83,10 @@ class TestEachScheduling:
class TestLoadScheduling: class TestLoadScheduling:
def test_schedule_load_simple(self): def test_schedule_load_simple(self):
node1 = MockNode()
node2 = MockNode()
sched = LoadScheduling(2) sched = LoadScheduling(2)
sched.addnode(node1) sched.addnode(MockNode())
sched.addnode(node2) sched.addnode(MockNode())
node1, node2 = sched.nodes
collection = ["a.py::test_1", "a.py::test_2"] collection = ["a.py::test_1", "a.py::test_2"]
assert not sched.collection_is_completed assert not sched.collection_is_completed
sched.addnode_collection(node1, collection) sched.addnode_collection(node1, collection)
@@ -98,38 +96,40 @@ class TestLoadScheduling:
assert sched.node2collection[node1] == collection assert sched.node2collection[node1] == collection
assert sched.node2collection[node2] == collection assert sched.node2collection[node2] == collection
sched.init_distribute() sched.init_distribute()
assert sched.tests_finished()
assert len(node1.sent) == 1
assert len(node2.sent) == 1
x = sorted(node1.sent + node2.sent)
assert x == [0, 1]
sched.remove_item(node1, node1.sent[0])
sched.remove_item(node2, node2.sent[0])
assert sched.tests_finished()
assert not sched.pending assert not sched.pending
assert not sched.tests_finished()
assert len(node1.sent) == 2
assert len(node2.sent) == 0
assert node1.sent == [0, 1]
sched.remove_item(node1, node1.sent[0])
assert sched.tests_finished()
sched.remove_item(node1, node1.sent[1])
assert sched.tests_finished()
def test_init_distribute_chunksize(self): def test_init_distribute_chunksize(self):
sched = LoadScheduling(2) sched = LoadScheduling(2)
node1 = MockNode() sched.addnode(MockNode())
node2 = MockNode() sched.addnode(MockNode())
sched.addnode(node1) node1, node2 = sched.nodes
sched.addnode(node2) col = ["xyz"] * (6)
col = ["xyz"] * (3)
sched.addnode_collection(node1, col) sched.addnode_collection(node1, col)
sched.addnode_collection(node2, col) sched.addnode_collection(node2, col)
sched.init_distribute() sched.init_distribute()
#assert not sched.tests_finished() #assert not sched.tests_finished()
sent1 = node1.sent sent1 = node1.sent
sent2 = node2.sent sent2 = node2.sent
chunkitems = col[:1] assert sent1 == [0, 1]
assert (sent1 == [0] and sent2 == [1]) or ( assert sent2 == [2, 3]
sent1 == [1] and sent2 == [0]) assert sched.pending == [4, 5]
assert sched.node2pending[node1] == sent1 assert sched.node2pending[node1] == sent1
assert sched.node2pending[node2] == sent2 assert sched.node2pending[node2] == sent2
assert len(sched.pending) == 1 assert len(sched.pending) == 2
for node in (node1, node2): sched.remove_item(node1, 0)
for i in sched.node2pending[node]: assert node1.sent == [0, 1, 4]
sched.remove_item(node, i) assert sched.pending == [5]
assert node2.sent == [2, 3]
sched.remove_item(node1, 1)
assert node1.sent == [0, 1, 4, 5]
assert not sched.pending assert not sched.pending
def test_add_remove_node(self): def test_add_remove_node(self):
@@ -144,6 +144,37 @@ class TestLoadScheduling:
crashitem = sched.remove_node(node) crashitem = sched.remove_node(node)
assert crashitem == collection[0] assert crashitem == collection[0]
def test_different_tests_collected(self, testdir):
"""
Test that LoadScheduling is reporting collection errors when
different test ids are collected by slaves.
"""
class CollectHook(object):
"""
Dummy hook that stores collection reports.
"""
def __init__(self):
self.reports = []
def pytest_collectreport(self, report):
self.reports.append(report)
collect_hook = CollectHook()
config = testdir.parseconfig()
config.pluginmanager.register(collect_hook, "collect_hook")
node1 = MockNode()
node2 = MockNode()
sched = LoadScheduling(2, config=config)
sched.addnode(node1)
sched.addnode(node2)
sched.addnode_collection(node1, ["a.py::test_1"])
sched.addnode_collection(node2, ["a.py::test_2"])
sched.init_distribute()
assert len(collect_hook.reports) == 1
rep = collect_hook.reports[0]
assert 'Different tests were collected between' in rep.longrepr
class TestDistReporter: class TestDistReporter:
@@ -179,7 +210,7 @@ class TestDistReporter:
def test_report_collection_diff_equal(): def test_report_collection_diff_equal():
"""Test reporting of equal collections.""" """Test reporting of equal collections."""
from_collection = to_collection = ['aaa', 'bbb', 'ccc'] from_collection = to_collection = ['aaa', 'bbb', 'ccc']
assert report_collection_diff(from_collection, to_collection, 1, 2) assert report_collection_diff(from_collection, to_collection, 1, 2) is None
def test_report_collection_diff_different(): def test_report_collection_diff_different():
@@ -202,10 +233,8 @@ def test_report_collection_diff_different():
'-YYY' '-YYY'
) )
try: msg = report_collection_diff(from_collection, to_collection, 1, 2)
report_collection_diff(from_collection, to_collection, 1, 2) assert msg == error_message
except AssertionError as e:
assert py.builtin._totext(e) == error_message
@pytest.mark.xfail(reason="duplicate test ids not supported yet") @pytest.mark.xfail(reason="duplicate test ids not supported yet")
def test_pytest_issue419(testdir): def test_pytest_issue419(testdir):

View File

@@ -42,7 +42,7 @@ class TestStatRecorder:
def test_dirchange(self, tmpdir): def test_dirchange(self, tmpdir):
tmp = tmpdir tmp = tmpdir
hello = tmp.ensure("dir", "hello.py") tmp.ensure("dir", "hello.py")
sd = StatRecorder([tmp]) sd = StatRecorder([tmp])
assert not sd.fil(tmp.join("dir")) assert not sd.fil(tmp.join("dir"))

View File

@@ -13,7 +13,7 @@ def test_dist_incompatibility_messages(testdir):
assert "incompatible" in result.stderr.str() assert "incompatible" in result.stderr.str()
def test_dist_options(testdir): def test_dist_options(testdir):
from xdist.plugin import check_options from xdist.plugin import pytest_cmdline_main as check_options
config = testdir.parseconfigure("-n 2") config = testdir.parseconfigure("-n 2")
check_options(config) check_options(config)
assert config.option.dist == "load" assert config.option.dist == "load"
@@ -67,4 +67,3 @@ class TestDistOptions:
assert py.path.local('y') in roots assert py.path.local('y') in roots
assert py.path.local('z') in roots assert py.path.local('z') in roots
assert testdir.tmpdir.join('x') in roots assert testdir.tmpdir.join('x') in roots

View File

@@ -3,7 +3,6 @@ from xdist.slavemanage import SlaveController, unserialize_report
from xdist.remote import serialize_report from xdist.remote import serialize_report
import execnet import execnet
queue = py.builtin._tryimport("queue", "Queue") queue = py.builtin._tryimport("queue", "Queue")
from py.builtin import print_
import marshal import marshal
WAIT_TIMEOUT = 10.0 WAIT_TIMEOUT = 10.0
@@ -26,7 +25,7 @@ class SlaveSetup:
use_callback = False use_callback = False
def __init__(self, request): def __init__(self, request):
self.testdir = testdir = request.getfuncargvalue("testdir") self.testdir = request.getfuncargvalue("testdir")
self.request = request self.request = request
self.events = queue.Queue() self.events = queue.Queue()
@@ -140,7 +139,7 @@ class TestReportSerialization:
class TestSlaveInteractor: class TestSlaveInteractor:
def test_basic_collect_and_runtests(self, slave): def test_basic_collect_and_runtests(self, slave):
p = slave.testdir.makepyfile(""" slave.testdir.makepyfile("""
def test_func(): def test_func():
pass pass
""") """)
@@ -170,7 +169,7 @@ class TestSlaveInteractor:
assert 'slaveoutput' in ev.kwargs assert 'slaveoutput' in ev.kwargs
def test_remote_collect_skip(self, slave): def test_remote_collect_skip(self, slave):
p = slave.testdir.makepyfile(""" slave.testdir.makepyfile("""
import py import py
py.test.skip("hello") py.test.skip("hello")
""") """)
@@ -187,7 +186,7 @@ class TestSlaveInteractor:
assert not ev.kwargs['ids'] assert not ev.kwargs['ids']
def test_remote_collect_fail(self, slave): def test_remote_collect_fail(self, slave):
p = slave.testdir.makepyfile("""aasd qwe""") slave.testdir.makepyfile("""aasd qwe""")
slave.setup() slave.setup()
ev = slave.popevent("collectionstart") ev = slave.popevent("collectionstart")
assert not ev.kwargs assert not ev.kwargs
@@ -201,7 +200,7 @@ class TestSlaveInteractor:
assert not ev.kwargs['ids'] assert not ev.kwargs['ids']
def test_runtests_all(self, slave): def test_runtests_all(self, slave):
p = slave.testdir.makepyfile(""" slave.testdir.makepyfile("""
def test_func(): pass def test_func(): pass
def test_func2(): pass def test_func2(): pass
""") """)
@@ -228,7 +227,7 @@ class TestSlaveInteractor:
def test_happy_run_events_converted(self, testdir, slave): def test_happy_run_events_converted(self, testdir, slave):
py.test.xfail("implement a simple test for event production") py.test.xfail("implement a simple test for event production")
assert not slave.use_callback assert not slave.use_callback
p = slave.testdir.makepyfile(""" slave.testdir.makepyfile("""
def test_func(): def test_func():
pass pass
""") """)

View File

@@ -1,28 +1,35 @@
import py import py
import os import pytest
import execnet import execnet
from _pytest.pytester import HookRecorder
from xdist import slavemanage, newhooks
from xdist.slavemanage import HostRSync, NodeManager from xdist.slavemanage import HostRSync, NodeManager
pytest_plugins = "pytester", pytest_plugins = "pytester"
def pytest_funcarg__hookrecorder(request): def pytest_funcarg__hookrecorder(request, config):
_pytest = request.getfuncargvalue('_pytest') hookrecorder = HookRecorder(config.pluginmanager)
config = request.getfuncargvalue('config') if hasattr(hookrecorder, "start_recording"):
return _pytest.gethookrecorder(config.hook) hookrecorder.start_recording(newhooks)
request.addfinalizer(hookrecorder.finish_recording)
return hookrecorder
def pytest_funcarg__config(request): def pytest_funcarg__config(testdir):
testdir = request.getfuncargvalue("testdir") return testdir.parseconfig()
config = testdir.parseconfig()
return config
def pytest_funcarg__mysetup(request): def pytest_funcarg__mysetup(tmpdir):
class mysetup: class mysetup:
def __init__(self, request): source = tmpdir.mkdir("source")
temp = request.getfuncargvalue("tmpdir") dest = tmpdir.mkdir("dest")
self.source = temp.mkdir("source") return mysetup()
self.dest = temp.mkdir("dest")
request.getfuncargvalue("_pytest") @pytest.fixture
return mysetup(request) def slavecontroller(monkeypatch):
class MockController(object):
def __init__(self, *args): pass
def setup(self): pass
monkeypatch.setattr(slavemanage, 'SlaveController', MockController)
return MockController
class TestNodeManagerPopen: class TestNodeManagerPopen:
def test_popen_no_default_chdir(self, config): def test_popen_no_default_chdir(self, config):
@@ -36,9 +43,9 @@ class TestNodeManagerPopen:
for spec in NodeManager(config, l, defaultchdir="abc").specs: for spec in NodeManager(config, l, defaultchdir="abc").specs:
assert spec.chdir == "abc" assert spec.chdir == "abc"
def test_popen_makegateway_events(self, config, hookrecorder, _pytest): def test_popen_makegateway_events(self, config, hookrecorder, slavecontroller):
hm = NodeManager(config, ["popen"] * 2) hm = NodeManager(config, ["popen"] * 2)
hm.makegateways() hm.setup_nodes(None)
call = hookrecorder.popcall("pytest_xdist_setupnodes") call = hookrecorder.popcall("pytest_xdist_setupnodes")
assert len(call.specs) == 2 assert len(call.specs) == 2
@@ -51,10 +58,10 @@ class TestNodeManagerPopen:
hm.teardown_nodes() hm.teardown_nodes()
assert not len(hm.group) assert not len(hm.group)
def test_popens_rsync(self, config, mysetup): def test_popens_rsync(self, config, mysetup, slavecontroller):
source = mysetup.source source = mysetup.source
hm = NodeManager(config, ["popen"] * 2) hm = NodeManager(config, ["popen"] * 2)
hm.makegateways() hm.setup_nodes(None)
assert len(hm.group) == 2 assert len(hm.group) == 2
for gw in hm.group: for gw in hm.group:
class pseudoexec: class pseudoexec:
@@ -65,19 +72,21 @@ class TestNodeManagerPopen:
pass pass
gw.remote_exec = pseudoexec gw.remote_exec = pseudoexec
l = [] l = []
hm.rsync(source, notify=lambda *args: l.append(args)) for gw in hm.group:
hm.rsync(gw, source, notify=lambda *args: l.append(args))
assert not l assert not l
hm.teardown_nodes() hm.teardown_nodes()
assert not len(hm.group) assert not len(hm.group)
assert "sys.path.insert" in gw.remote_exec.args[0] assert "sys.path.insert" in gw.remote_exec.args[0]
def test_rsync_popen_with_path(self, config, mysetup): def test_rsync_popen_with_path(self, config, mysetup, slavecontroller):
source, dest = mysetup.source, mysetup.dest source, dest = mysetup.source, mysetup.dest
hm = NodeManager(config, ["popen//chdir=%s" %dest] * 1) hm = NodeManager(config, ["popen//chdir=%s" % dest] * 1)
hm.makegateways() hm.setup_nodes(None)
source.ensure("dir1", "dir2", "hello") source.ensure("dir1", "dir2", "hello")
l = [] l = []
hm.rsync(source, notify=lambda *args: l.append(args)) for gw in hm.group:
hm.rsync(gw, source, notify=lambda *args: l.append(args))
assert len(l) == 1 assert len(l) == 1
assert l[0] == ("rsyncrootready", hm.group['gw0'].spec, source) assert l[0] == ("rsyncrootready", hm.group['gw0'].spec, source)
hm.teardown_nodes() hm.teardown_nodes()
@@ -86,12 +95,15 @@ class TestNodeManagerPopen:
assert dest.join("dir1", "dir2").check() assert dest.join("dir1", "dir2").check()
assert dest.join("dir1", "dir2", 'hello').check() assert dest.join("dir1", "dir2", 'hello').check()
def test_rsync_same_popen_twice(self, config, mysetup, hookrecorder): def test_rsync_same_popen_twice(self, config, mysetup,
hookrecorder, slavecontroller):
source, dest = mysetup.source, mysetup.dest source, dest = mysetup.source, mysetup.dest
hm = NodeManager(config, ["popen//chdir=%s" %dest] * 2) hm = NodeManager(config, ["popen//chdir=%s" % dest] * 2)
hm.makegateways() hm.roots = []
hm.setup_nodes(None)
source.ensure("dir1", "dir2", "hello") source.ensure("dir1", "dir2", "hello")
hm.rsync(source) gw = hm.group[0]
hm.rsync(gw, source)
call = hookrecorder.popcall("pytest_xdist_rsyncstart") call = hookrecorder.popcall("pytest_xdist_rsyncstart")
assert call.source == source assert call.source == source
assert len(call.gateways) == 1 assert len(call.gateways) == 1
@@ -99,16 +111,8 @@ class TestNodeManagerPopen:
call = hookrecorder.popcall("pytest_xdist_rsyncfinish") call = hookrecorder.popcall("pytest_xdist_rsyncfinish")
class TestHRSync: class TestHRSync:
def pytest_funcarg__mysetup(self, request):
class mysetup:
def __init__(self, request):
tmp = request.getfuncargvalue('tmpdir')
self.source = tmp.mkdir("source")
self.dest = tmp.mkdir("dest")
return mysetup(request)
def test_hrsync_filter(self, mysetup): def test_hrsync_filter(self, mysetup):
source, dest = mysetup.source, mysetup.dest source, _ = mysetup.source, mysetup.dest # noqa
source.ensure("dir", "file.txt") source.ensure("dir", "file.txt")
source.ensure(".svn", "entries") source.ensure(".svn", "entries")
source.ensure(".somedotfile", "moreentries") source.ensure(".somedotfile", "moreentries")
@@ -136,10 +140,10 @@ class TestHRSync:
class TestNodeManager: class TestNodeManager:
@py.test.mark.xfail @py.test.mark.xfail(run=False)
def test_rsync_roots_no_roots(self, testdir, mysetup): def test_rsync_roots_no_roots(self, testdir, mysetup):
mysetup.source.ensure("dir1", "file1").write("hello") mysetup.source.ensure("dir1", "file1").write("hello")
config = testdir.parseconfig(source) config = testdir.parseconfig(mysetup.source)
nodemanager = NodeManager(config, ["popen//chdir=%s" % mysetup.dest]) nodemanager = NodeManager(config, ["popen//chdir=%s" % mysetup.dest])
#assert nodemanager.config.topdir == source == config.topdir #assert nodemanager.config.topdir == source == config.topdir
nodemanager.makegateways() nodemanager.makegateways()
@@ -152,7 +156,7 @@ class TestNodeManager:
assert p.join("dir1").check() assert p.join("dir1").check()
assert p.join("dir1", "file1").check() assert p.join("dir1", "file1").check()
def test_popen_rsync_subdir(self, testdir, mysetup): def test_popen_rsync_subdir(self, testdir, mysetup, slavecontroller):
source, dest = mysetup.source, mysetup.dest source, dest = mysetup.source, mysetup.dest
dir1 = mysetup.source.mkdir("dir1") dir1 = mysetup.source.mkdir("dir1")
dir2 = dir1.mkdir("dir2") dir2 = dir1.mkdir("dir2")
@@ -164,8 +168,7 @@ class TestNodeManager:
"--rsyncdir", rsyncroot, "--rsyncdir", rsyncroot,
source, source,
)) ))
nodemanager.makegateways() nodemanager.setup_nodes(None) # calls .rsync_roots()
nodemanager.rsync_roots()
if rsyncroot == source: if rsyncroot == source:
dest = dest.join("source") dest = dest.join("source")
assert dest.join("dir1").check() assert dest.join("dir1").check()
@@ -173,7 +176,7 @@ class TestNodeManager:
assert dest.join("dir1", "dir2", 'hello').check() assert dest.join("dir1", "dir2", 'hello').check()
nodemanager.teardown_nodes() nodemanager.teardown_nodes()
def test_init_rsync_roots(self, testdir, mysetup): def test_init_rsync_roots(self, testdir, mysetup, slavecontroller):
source, dest = mysetup.source, mysetup.dest source, dest = mysetup.source, mysetup.dest
dir2 = source.ensure("dir1", "dir2", dir=1) dir2 = source.ensure("dir1", "dir2", dir=1)
source.ensure("dir1", "somefile", dir=1) source.ensure("dir1", "somefile", dir=1)
@@ -185,20 +188,19 @@ class TestNodeManager:
""")) """))
config = testdir.parseconfig(source) config = testdir.parseconfig(source)
nodemanager = NodeManager(config, ["popen//chdir=%s" % dest]) nodemanager = NodeManager(config, ["popen//chdir=%s" % dest])
nodemanager.makegateways() nodemanager.setup_nodes(None) # calls .rsync_roots()
nodemanager.rsync_roots()
assert dest.join("dir2").check() assert dest.join("dir2").check()
assert not dest.join("dir1").check() assert not dest.join("dir1").check()
assert not dest.join("bogus").check() assert not dest.join("bogus").check()
def test_rsyncignore(self, testdir, mysetup): def test_rsyncignore(self, testdir, mysetup, slavecontroller):
source, dest = mysetup.source, mysetup.dest source, dest = mysetup.source, mysetup.dest
dir2 = source.ensure("dir1", "dir2", dir=1) dir2 = source.ensure("dir1", "dir2", dir=1)
dir5 = source.ensure("dir5", "dir6", "bogus") source.ensure("dir5", "dir6", "bogus")
dirf = source.ensure("dir5", "file") source.ensure("dir5", "file")
dir2.ensure("hello") dir2.ensure("hello")
dirfoo = source.ensure("foo", "bar") source.ensure("foo", "bar")
dirbar = source.ensure("bar", "foo") source.ensure("bar", "foo")
source.join("tox.ini").write(py.std.textwrap.dedent(""" source.join("tox.ini").write(py.std.textwrap.dedent("""
[pytest] [pytest]
rsyncdirs = dir1 dir5 rsyncdirs = dir1 dir5
@@ -207,24 +209,22 @@ class TestNodeManager:
config = testdir.parseconfig(source) config = testdir.parseconfig(source)
config.option.rsyncignore = ['bar'] config.option.rsyncignore = ['bar']
nodemanager = NodeManager(config, ["popen//chdir=%s" % dest]) nodemanager = NodeManager(config, ["popen//chdir=%s" % dest])
nodemanager.makegateways() nodemanager.setup_nodes(None) # calls .rsync_roots()
nodemanager.rsync_roots()
assert dest.join("dir1").check() assert dest.join("dir1").check()
assert not dest.join("dir1", "dir2").check() assert not dest.join("dir1", "dir2").check()
assert dest.join("dir5","file").check() assert dest.join("dir5", "file").check()
assert not dest.join("dir6").check() assert not dest.join("dir6").check()
assert not dest.join('foo').check() assert not dest.join('foo').check()
assert not dest.join('bar').check() assert not dest.join('bar').check()
def test_optimise_popen(self, testdir, mysetup): def test_optimise_popen(self, testdir, mysetup, slavecontroller):
source, dest = mysetup.source, mysetup.dest source = mysetup.source
specs = ["popen"] * 3 specs = ["popen"] * 3
source.join("conftest.py").write("rsyncdirs = ['a']") source.join("conftest.py").write("rsyncdirs = ['a']")
source.ensure('a', dir=1) source.ensure('a', dir=1)
config = testdir.parseconfig(source) config = testdir.parseconfig(source)
nodemanager = NodeManager(config, specs) nodemanager = NodeManager(config, specs)
nodemanager.makegateways() nodemanager.setup_nodes(None) # calls .rysnc_roots()
nodemanager.rsync_roots()
for gwspec in nodemanager.specs: for gwspec in nodemanager.specs:
assert gwspec._samefilesystem() assert gwspec._samefilesystem()
assert not gwspec.chdir assert not gwspec.chdir
@@ -235,8 +235,6 @@ class TestNodeManager:
pass pass
""") """)
reprec = testdir.inline_run("-d", "--rsyncdir=%s" % testdir.tmpdir, reprec = testdir.inline_run("-d", "--rsyncdir=%s" % testdir.tmpdir,
"--tx", specssh, testdir.tmpdir) "--tx", specssh, testdir.tmpdir)
rep, = reprec.getreports("pytest_runtest_logreport") rep, = reprec.getreports("pytest_runtest_logreport")
assert rep.passed assert rep.passed

35
tox.ini
View File

@@ -1,26 +1,27 @@
[tox] [tox]
envlist=py26,py32,py33,py27,py27-pexpect,py33-pexpect,py26,py26-old,py33-old envlist=
py{26,33,34,27}-pytest2{4,5,6,7},py{27,34}-pytest27-pexpect,flakes
[testenv] [testenv]
changedir=testing changedir=testing
deps=pytest>=2.5.1 deps =
commands= py.test --junitxml={envlogdir}/junit-{envname}.xml [] pycmd
pytest24: pytest~=2.4.0
pytest25: pytest~=2.5.0
[testenv:py27-pexpect] pytest26: pytest~=2.6.1
deps={[testenv]deps} pytest27: pytest~=2.7.2
pexpect pexpect: pexpect
[testenv:py33-pexpect] commands=
deps={[testenv]deps} # always clean to avoid code unmarshal mismatch on old python/pytest
pexpect py.cleanup -aq
py.test {posargs}
[testenv:py26-old] [testenv:flakes]
deps= changedir=
pytest==2.4.2 deps = pytest-flakes>=0.2
commands = py.test --flakes -m flakes testing xdist
[testenv:py33-old]
basepython = python3.3
deps=
pytest==2.4.2
[pytest] [pytest]
addopts = -rsfxX addopts = -rsfxX

View File

@@ -1,2 +1,2 @@
# __all__ = ['__version__']
__version__ = '1.10' from xdist._version import version as __version__

56
xdist/boxed.py Normal file
View File

@@ -0,0 +1,56 @@
import py
def pytest_addoption(parser):
group = parser.getgroup("xdist", "distributed and subprocess testing")
group.addoption('--boxed',
action="store_true", dest="boxed", default=False,
help="box each test run in a separate process (unix)")
def pytest_runtest_protocol(item):
if item.config.getvalue("boxed"):
reports = forked_run_report(item)
for rep in reports:
item.ihook.pytest_runtest_logreport(report=rep)
return True
def forked_run_report(item):
# for now, we run setup/teardown in the subprocess
# XXX optionally allow sharing of setup/teardown
from _pytest.runner import runtestprotocol
EXITSTATUS_TESTEXIT = 4
import marshal
from xdist.remote import serialize_report
from xdist.slavemanage import unserialize_report
def runforked():
try:
reports = runtestprotocol(item, log=False)
except KeyboardInterrupt:
py.std.os._exit(EXITSTATUS_TESTEXIT)
return marshal.dumps([serialize_report(x) for x in reports])
ff = py.process.ForkedFunc(runforked)
result = ff.waitfinish()
if result.retval is not None:
report_dumps = marshal.loads(result.retval)
return [unserialize_report("testreport", x) for x in report_dumps]
else:
if result.exitstatus == EXITSTATUS_TESTEXIT:
py.test.exit("forked test item %s raised Exit" %(item,))
return [report_process_crash(item, result)]
def report_process_crash(item, result):
path, lineno = item._getfslineno()
info = ("%s:%s: running the test CRASHED with signal %d" %
(path, lineno, result.signal))
from _pytest import runner
call = runner.CallInfo(lambda: 0/0, "???")
call.excinfo = info
rep = runner.pytest_runtest_makereport(item, call)
if result.out:
rep.sections.append(("captured stdout", result.out))
if result.err:
rep.sections.append(("captured stderr", result.err))
return rep

View File

@@ -1,5 +1,5 @@
import sys
import difflib import difflib
from _pytest.runner import CollectReport
import pytest import pytest
import py import py
@@ -10,35 +10,96 @@ queue = py.builtin._tryimport('queue', 'Queue')
class EachScheduling: class EachScheduling:
"""Implement scheduling of test items on all nodes
If a node gets added after the test run is started then it is
assumed to replace a node which got removed before it finished
it's collection. In this case it will only be used if a a node
with the same spec got removed earlier.
Any nodes added after the run is started will only get items
assigned if a node with matching spec was removed before it
finished all it's pending items. The new node will then be
assigned the remaining items from the removed node.
"""
def __init__(self, numnodes, log=None): def __init__(self, numnodes, log=None):
self.numnodes = numnodes self.numnodes = numnodes
self.node2collection = {} self.node2collection = {}
self.node2pending = {} self.node2pending = {}
self._started = []
self._removed2pending = {}
if log is None: if log is None:
self.log = py.log.Producer("eachsched") self.log = py.log.Producer("eachsched")
else: else:
self.log = log.loadsched self.log = log.eachsched
self.collection_is_completed = False self.collection_is_completed = False
@property
def nodes(self):
"""A list of all nodes in the scheduler."""
return list(self.node2pending.keys())
def hasnodes(self): def hasnodes(self):
return bool(self.node2pending) return bool(self.node2pending)
def haspending(self):
"""Return True if there are pending test items
This indicates that collection has finished and nodes are
still processing test items, so can be thought of as "the
scheduler is active".
"""
for pending in self.node2pending.values():
if pending:
return True
return False
def addnode(self, node): def addnode(self, node):
self.node2collection[node] = None assert node not in self.node2pending
self.node2pending[node] = []
def tests_finished(self): def tests_finished(self):
if not self.collection_is_completed: if not self.collection_is_completed:
return False return False
if self._removed2pending:
return False
for pending in self.node2pending.values():
if len(pending) >= 2:
return False
return True return True
def addnode_collection(self, node, collection): def addnode_collection(self, node, collection):
assert not self.collection_is_completed """Add the collected test items from a node
assert self.node2collection[node] is None
self.node2collection[node] = list(collection) Collection is complete once all nodes have submitted their
self.node2pending[node] = [] collection. In this case it's peding list is set to an empty
if len(self.node2pending) >= self.numnodes: list. When the collection is already completed this
self.collection_is_completed = True submission is from a node which was restarted to replace a
dead node. In this case we already assing the pending items
here. In either case ``.init_distribute()`` will instruct the
node to start running the required tests.
"""
assert node in self.node2pending
if not self.collection_is_completed:
self.node2collection[node] = list(collection)
self.node2pending[node] = []
if len(self.node2collection) >= self.numnodes:
self.collection_is_completed = True
elif self._removed2pending:
for deadnode in self._removed2pending:
if deadnode.gateway.spec == node.gateway.spec:
dead_collection = self.node2collection[deadnode]
if collection != dead_collection:
msg = report_collection_diff(dead_collection,
collection,
deadnode.gateway.id,
node.gateway.id)
self.log(msg)
return
pending = self._removed2pending.pop(deadnode)
self.node2pending[node] = pending
break
def remove_item(self, node, item_index, duration=0): def remove_item(self, node, item_index, duration=0):
self.node2pending[node].remove(item_index) self.node2pending[node].remove(item_index)
@@ -49,129 +110,309 @@ class EachScheduling:
if not pending: if not pending:
return return
crashitem = self.node2collection[node][pending.pop(0)] crashitem = self.node2collection[node][pending.pop(0)]
# XXX what about the rest of pending? if pending:
self._removed2pending[node] = pending
return crashitem return crashitem
def init_distribute(self): def init_distribute(self):
"""Schedule the test items on the nodes
If the node's pending list is empty it is a new node which
needs to run all the tests. If the pending list is already
populated (by ``.addnode_collection()``) then it replaces a
died node and we only need to run those tests.
"""
assert self.collection_is_completed assert self.collection_is_completed
for node, pending in self.node2pending.items(): for node, pending in self.node2pending.items():
node.send_runtest_all() if node in self._started:
pending[:] = range(len(self.node2collection[node])) continue
if not pending:
pending[:] = range(len(self.node2collection[node]))
node.send_runtest_all()
else:
node.send_runtest_some(pending)
self._started.append(node)
class LoadScheduling: class LoadScheduling:
def __init__(self, numnodes, log=None): """Implement load scheduling accross nodes.
This distributes the tests collected across all nodes so each test
is run just once. All nodes collect and submit the test suit and
when all collections are received it is verified they are
identical collections. Then the collection gets devided up in
chunks and chunks get submitted to nodes. Whenver a node finishes
an item they call ``.remove_item()`` which will trigger the
scheduler to assign more tests if the number of pending tests for
the node falls below a low-watermark.
When created ``numnodes`` defines how many nodes are expected to
submit a collection, this is used to know when all nodes have
finished collection or how large the chunks need to be created.
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 died 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 ``.init_distribute()`` 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.
"""
def __init__(self, numnodes, log=None, config=None):
self.numnodes = numnodes self.numnodes = numnodes
self.node2pending = {}
self.node2collection = {} self.node2collection = {}
self.node2pending = {}
self.pending = [] self.pending = []
self.collection = None
if log is None: if log is None:
self.log = py.log.Producer("loadsched") self.log = py.log.Producer("loadsched")
else: else:
self.log = log.loadsched self.log = log.loadsched
self.collection_is_completed = False self.config = config
@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
def haspending(self):
"""Return True if there are pending test items
This indicates that collection has finished and nodes are
still processing test items, so 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 hasnodes(self): def hasnodes(self):
"""Return True if nodes exist in the scheduler."""
return bool(self.node2pending) return bool(self.node2pending)
def addnode(self, node): def addnode(self, node):
"""Add a new node in the scheduler.
From now on the node will be allocated chunks of tests to
execute.
Called by the ``DSession.slave_slaveready`` hook when it
sucessfully bootstrapped a new node.
"""
assert node not in self.node2pending
self.node2pending[node] = [] self.node2pending[node] = []
def tests_finished(self): def tests_finished(self):
if not self.collection_is_completed or self.pending: """Return True if all tests have been executed by the nodes."""
if not self.collection_is_completed:
return False return False
#for items in self.node2pending.values(): if self.pending:
# if items: return False
# return False for pending in self.node2pending.values():
if len(pending) >= 2:
return False
return True return True
def addnode_collection(self, node, collection): def addnode_collection(self, node, collection):
assert not self.collection_is_completed """Add the collected test items from a node
The collection is stored in the ``.node2collection`` map.
Called by the ``DSession.slave_collectionfinish`` hook.
"""
assert node in self.node2pending assert node in self.node2pending
if self.collection_is_completed:
# A new node has been added later, perhaps an original one died.
assert self.collection # .init_distribute() should have
# been called by now
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) self.node2collection[node] = list(collection)
if len(self.node2collection) >= self.numnodes:
self.collection_is_completed = True
def remove_item(self, node, item_index, duration=0): def remove_item(self, node, item_index, duration=0):
node_pending = self.node2pending[node] """Mark test item as completed by node
node_pending.remove(item_index)
# pre-load items-to-test if the node may become ready
The duration it took to execute the item is used as a hint to
the scheduler.
This is called by the ``DSession.slave_testreport`` hook.
"""
self.node2pending[node].remove(item_index)
self.check_schedule(node, duration=duration)
def check_schedule(self, node, duration=0):
"""Maybe schedule new items on the node
If there are any globally pending nodes left then this will
check if the given node should be given any more tests. The
``duration`` of the last test is optionally used as a
heuristic to influence how many tests the node is assigned.
"""
if self.pending: if self.pending:
if duration >= 0.1 and node_pending: # how many nodes do we have?
# seems the node is doing long-running tests
# so let's rather wait with sending new items
return
# how many nodes do we have remaining per node roughly?
num_nodes = len(self.node2pending) num_nodes = len(self.node2pending)
# if our node goes below a heuristic minimum, fill it out to # if our node goes below a heuristic minimum, fill it out to
# heuristic maximum # heuristic maximum
items_per_node_min = max( items_per_node_min = max(2, len(self.pending) // num_nodes // 4)
1, len(self.pending) // num_nodes // 4) items_per_node_max = max(2, len(self.pending) // num_nodes // 2)
items_per_node_max = max( node_pending = self.node2pending[node]
1, len(self.pending) // num_nodes // 2) if len(node_pending) < items_per_node_min:
if len(node_pending) <= items_per_node_min: if duration >= 0.1 and len(node_pending) >= 2:
num_send = items_per_node_max - len(node_pending) + 1 # seems the node is doing long-running tests
# and has enough items to continue
# so let's rather wait with sending new items
return
num_send = items_per_node_max - len(node_pending)
self._send_tests(node, num_send) self._send_tests(node, num_send)
self.log("num items waiting for node:", len(self.pending)) self.log("num items waiting for node:", len(self.pending))
#self.log("node2pending:", self.node2pending)
def remove_node(self, node): def remove_node(self, node):
"""Remove an 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.slave_slavefinished`` and
``DSession.slave_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) pending = self.node2pending.pop(node)
if not pending: if not pending:
return return
# the node must have crashed on the item if there are pending ones
# The node crashed, reassing pending items
crashitem = self.collection[pending.pop(0)] crashitem = self.collection[pending.pop(0)]
self.pending.extend(pending) self.pending.extend(pending)
for node in self.node2pending:
self.check_schedule(node)
return crashitem return crashitem
def init_distribute(self): def init_distribute(self):
assert self.collection_is_completed """Initiate distribution of the test collection
# XXX allow nodes to have different collections
node_collection_items = list(self.node2collection.items())
first_node, col = node_collection_items[0]
for node, collection in node_collection_items[1:]:
report_collection_diff(
col,
collection,
first_node.gateway.id,
node.gateway.id,
)
# all collections are the same, good. Initiate scheduling of the items across the nodes. If this
# we now create an index gets called again later it behaves the same as calling
self.collection = col ``.check_schedule()`` on all nodes so that newly added nodes
self.pending[:] = range(len(col)) will start to be used.
if not col:
This is called by the ``DSession.slave_collectionfinish`` hook
if ``.collection_is_completed`` is True.
XXX Perhaps this method should have been called ".schedule()".
"""
assert self.collection_is_completed
# Initial distribution already happend, reschedule on all nodes
if self.collection is not None:
for node in self.nodes:
self.check_schedule(node)
return return
# XXX allow nodes to have different collections
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
# 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, 1) node_chunksize = max(items_per_node // 4, 2)
# and initialize each node with a chunk of tests # and initialize each node with a chunk of tests
for node in self.node2pending: for node in self.nodes:
self._send_tests(node, node_chunksize) self._send_tests(node, node_chunksize)
#f = open("/tmp/sent", "w")
def _send_tests(self, node, num): def _send_tests(self, node, num):
tests_per_node = self.pending[:num] tests_per_node = self.pending[:num]
#print >>self.f, "sent", node, tests_per_node
if tests_per_node: if tests_per_node:
del self.pending[:num] del self.pending[:num]
self.node2pending[node].extend(tests_per_node) self.node2pending[node].extend(tests_per_node)
node.send_runtest_some(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
def report_collection_diff(from_collection, to_collection, from_id, to_id): def report_collection_diff(from_collection, to_collection, from_id, to_id):
"""Report the collected test difference between two nodes. """Report the collected test difference between two nodes.
:returns: True if collections are equal. :returns: detailed message describing the difference between the given
collections, or None if they are equal.
:raises: AssertionError with a detailed error message describing the
difference between the collections.
""" """
if from_collection == to_collection: if from_collection == to_collection:
return True return None
diff = difflib.unified_diff( diff = difflib.unified_diff(
from_collection, from_collection,
@@ -185,13 +426,26 @@ def report_collection_diff(from_collection, to_collection, from_id, to_id):
'{diff}' '{diff}'
).format(from_id=from_id, to_id=to_id, diff='\n'.join(diff)) ).format(from_id=from_id, to_id=to_id, diff='\n'.join(diff))
msg = "\n".join([x.rstrip() for x in error_message.split("\n")]) msg = "\n".join([x.rstrip() for x in error_message.split("\n")])
raise AssertionError(msg) return msg
class Interrupted(KeyboardInterrupt): class Interrupted(KeyboardInterrupt):
""" signals an immediate interruption. """ """ signals an immediate interruption. """
class DSession: class DSession:
"""A py.test plugin which runs a distributed test session
At the beginning of the test session this creates a NodeManager
instance which creates and starts all nodes. Nodes then emit
events processed in the pytest_runtestloop hook using the slave_*
methods.
Once a node is started it will automatically start running the
py.test mainloop with some custom hooks. This means a node
automatically starts collecting tests. Once tests are collected
it will wait for instructions.
"""
def __init__(self, config): def __init__(self, config):
self.config = config self.config = config
self.log = py.log.Producer("dsession") self.log = py.log.Producer("dsession")
@@ -201,7 +455,13 @@ class DSession:
self.countfailures = 0 self.countfailures = 0
self.maxfail = config.getvalue("maxfail") self.maxfail = config.getvalue("maxfail")
self.queue = queue.Queue() self.queue = queue.Queue()
self._session = None
self._failed_collection_errors = {} self._failed_collection_errors = {}
self._active_nodes = set()
self._failed_nodes_count = 0
self._max_slave_restart = self.config.getoption('max_slave_restart')
if self._max_slave_restart is not None:
self._max_slave_restart = int(self._max_slave_restart)
try: try:
self.terminal = config.pluginmanager.getplugin("terminalreporter") self.terminal = config.pluginmanager.getplugin("terminalreporter")
except KeyError: except KeyError:
@@ -210,20 +470,37 @@ class DSession:
self.trdist = TerminalDistReporter(config) self.trdist = TerminalDistReporter(config)
config.pluginmanager.register(self.trdist, "terminaldistreporter") config.pluginmanager.register(self.trdist, "terminaldistreporter")
@property
def session_finished(self):
"""Return True if the distributed session has finished
This means all nodes have executed all test items. This is
used to by pytest_runtestloop to break out of it's loop.
"""
return bool(self.shuttingdown and not self._active_nodes)
def report_line(self, line): def report_line(self, line):
if self.terminal and self.config.option.verbose >= 0: if self.terminal and self.config.option.verbose >= 0:
self.terminal.write_line(line) self.terminal.write_line(line)
@pytest.mark.trylast @pytest.mark.trylast
def pytest_sessionstart(self, session): def pytest_sessionstart(self, session):
"""Creates and starts the nodes.
The nodes are setup to put their events onto self.queue. As
soon as nodes start they will emit the slave_slaveready event.
"""
self.nodemanager = NodeManager(self.config) self.nodemanager = NodeManager(self.config)
self.nodemanager.setup_nodes(putevent=self.queue.put) nodes = self.nodemanager.setup_nodes(putevent=self.queue.put)
self._active_nodes.update(nodes)
self._session = session
def pytest_sessionfinish(self, session): def pytest_sessionfinish(self, session):
""" teardown any resources after a test run. """ """Shutdown all nodes."""
nm = getattr(self, 'nodemanager', None) # if not fully initialized nm = getattr(self, 'nodemanager', None) # if not fully initialized
if nm is not None: if nm is not None:
nm.teardown_nodes() nm.teardown_nodes()
self._session = None
def pytest_collection(self): def pytest_collection(self):
# prohibit collection of test items in master process # prohibit collection of test items in master process
@@ -233,13 +510,13 @@ class DSession:
numnodes = len(self.nodemanager.specs) numnodes = len(self.nodemanager.specs)
dist = self.config.getvalue("dist") dist = self.config.getvalue("dist")
if dist == "load": if dist == "load":
self.sched = LoadScheduling(numnodes, log=self.log) self.sched = LoadScheduling(numnodes, log=self.log,
config=self.config)
elif dist == "each": elif dist == "each":
self.sched = EachScheduling(numnodes, log=self.log) self.sched = EachScheduling(numnodes, log=self.log)
else: else:
assert 0, dist assert 0, dist
self.shouldstop = False self.shouldstop = False
self.session_finished = False
while not self.session_finished: while not self.session_finished:
self.loop_once() self.loop_once()
if self.shouldstop: if self.shouldstop:
@@ -247,7 +524,7 @@ class DSession:
return True return True
def loop_once(self): def loop_once(self):
""" process one callback from one of the slaves. """ """Process one callback from one of the slaves."""
while 1: while 1:
try: try:
eventcall = self.queue.get(timeout=2.0) eventcall = self.queue.get(timeout=2.0)
@@ -268,26 +545,40 @@ class DSession:
# #
def slave_slaveready(self, node, slaveinfo): def slave_slaveready(self, node, slaveinfo):
"""Emitted when a node first starts up.
This adds the node to the scheduler, nodes continue with
collection without any further input.
"""
node.slaveinfo = slaveinfo node.slaveinfo = slaveinfo
node.slaveinfo['id'] = node.gateway.id node.slaveinfo['id'] = node.gateway.id
node.slaveinfo['spec'] = node.gateway.spec node.slaveinfo['spec'] = node.gateway.spec
self.config.hook.pytest_testnodeready(node=node) self.config.hook.pytest_testnodeready(node=node)
self.sched.addnode(node)
if self.shuttingdown: if self.shuttingdown:
node.shutdown() node.shutdown()
else:
self.sched.addnode(node)
def slave_slavefinished(self, node): def slave_slavefinished(self, node):
"""Emitted when node executes its pytest_sessionfinish hook.
Removes the node from the scheduler.
The node might not be the scheduler if it had not emitted
slaveready before shutdown was triggered.
"""
self.config.hook.pytest_testnodedown(node=node, error=None) self.config.hook.pytest_testnodedown(node=node, error=None)
if node.slaveoutput['exitstatus'] == 2: # keyboard-interrupt if node.slaveoutput['exitstatus'] == 2: # keyboard-interrupt
self.shouldstop = "%s received keyboard-interrupt" % (node,) self.shouldstop = "%s received keyboard-interrupt" % (node,)
self.slave_errordown(node, "keyboard-interrupt") self.slave_errordown(node, "keyboard-interrupt")
return return
crashitem = self.sched.remove_node(node) if node in self.sched.nodes:
#assert not crashitem, (crashitem, node) crashitem = self.sched.remove_node(node)
if self.shuttingdown and not self.sched.hasnodes(): assert not crashitem, (crashitem, node)
self.session_finished = True self._active_nodes.remove(node)
def slave_errordown(self, node, error): def slave_errordown(self, node, error):
"""Emitted by the SlaveController when a node dies."""
self.config.hook.pytest_testnodedown(node=node, error=error) self.config.hook.pytest_testnodedown(node=node, error=error)
try: try:
crashitem = self.sched.remove_node(node) crashitem = self.sched.remove_node(node)
@@ -296,41 +587,85 @@ class DSession:
else: else:
if crashitem: if crashitem:
self.handle_crashitem(crashitem, node) self.handle_crashitem(crashitem, node)
#self.report_line("item crashed on node: %s" % crashitem)
if not self.sched.hasnodes(): self._failed_nodes_count += 1
self.session_finished = True maximum_reached = (self._max_slave_restart is not None and
self._failed_nodes_count > self._max_slave_restart)
if maximum_reached:
if self._max_slave_restart == 0:
msg = 'Slave restarting disabled'
else:
msg = "Maximum crashed slaves reached: %d" % \
self._max_slave_restart
self.report_line(msg)
else:
self.report_line("Replacing crashed slave %s" % node.gateway.id)
self._clone_node(node)
self._active_nodes.remove(node)
def slave_collectionfinish(self, node, ids): def slave_collectionfinish(self, node, ids):
"""Slave has finished test collection.
This adds the collection for this node to the scheduler. If
the scheduler indicates collection is finished (i.e. all
initial nodes have submitted their collection), then tells the
scheduler to schedule the collected items. When initiating
scheduling the first time it logs which scheduler is in use.
"""
if self.shuttingdown:
return
# tell session which items were effectively collected otherwise
# the master node will finish the session with EXIT_NOTESTSCOLLECTED
self._session.testscollected = len(ids)
self.sched.addnode_collection(node, ids) self.sched.addnode_collection(node, ids)
if self.terminal: if self.terminal:
self.trdist.setstatus(node.gateway.spec, "[%d]" %(len(ids))) self.trdist.setstatus(node.gateway.spec, "[%d]" % (len(ids)))
if self.sched.collection_is_completed: if self.sched.collection_is_completed:
if self.terminal: if self.terminal and not self.sched.haspending():
self.trdist.ensure_show_status() self.trdist.ensure_show_status()
self.terminal.write_line("") self.terminal.write_line("")
self.terminal.write_line("scheduling tests via %s" %( self.terminal.write_line("scheduling tests via %s" % (
self.sched.__class__.__name__)) self.sched.__class__.__name__))
self.sched.init_distribute() self.sched.init_distribute()
def slave_logstart(self, node, nodeid, location): def slave_logstart(self, node, nodeid, location):
"""Emitted when a node calls the pytest_runtest_logstart hook."""
self.config.hook.pytest_runtest_logstart( self.config.hook.pytest_runtest_logstart(
nodeid=nodeid, location=location) nodeid=nodeid, location=location)
def slave_testreport(self, node, rep): def slave_testreport(self, node, rep):
if not (rep.passed and rep.when != "call"): """Emitted when a node calls the pytest_runtest_logreport hook.
if rep.when in ("setup", "call"):
self.sched.remove_item(node, rep.item_index, rep.duration) If the node indicates it is finished with a test item remove
the item from the pending list in the scheduler.
"""
if rep.when == "call" or (rep.when == "setup" and not rep.passed):
self.sched.remove_item(node, rep.item_index, rep.duration)
#self.report_line("testreport %s: %s" %(rep.id, rep.status)) #self.report_line("testreport %s: %s" %(rep.id, rep.status))
rep.node = node rep.node = node
self.config.hook.pytest_runtest_logreport(report=rep) self.config.hook.pytest_runtest_logreport(report=rep)
self._handlefailures(rep) self._handlefailures(rep)
def slave_collectreport(self, node, rep): def slave_collectreport(self, node, rep):
"""Emitted when a node calls the pytest_collectreport hook."""
if rep.failed: if rep.failed:
self._failed_slave_collectreport(node, rep) self._failed_slave_collectreport(node, rep)
def _clone_node(self, node):
"""Return new node based on an existing one.
This is normally for when a node died, this will copy the spec
of the existing node and create a new one with a new id. The
new node will have been setup so will start calling the
"slave_*" hooks and do work soon.
"""
spec = node.gateway.spec
spec.id = None
self.nodemanager.group.allocate_id(spec)
node = self.nodemanager.setup_node(spec, self.queue.put)
self._active_nodes.add(node)
return node
def _failed_slave_collectreport(self, node, rep): def _failed_slave_collectreport(self, node, rep):
# Check we haven't already seen this report (from # Check we haven't already seen this report (from
# another slave). # another slave).
@@ -349,19 +684,21 @@ class DSession:
def triggershutdown(self): def triggershutdown(self):
self.log("triggering shutdown") self.log("triggering shutdown")
self.shuttingdown = True self.shuttingdown = True
for node in self.sched.node2pending: for node in self.sched.nodes:
node.shutdown() node.shutdown()
def handle_crashitem(self, nodeid, slave): def handle_crashitem(self, nodeid, slave):
# XXX get more reporting info by recording pytest_runtest_logstart? # XXX get more reporting info by recording pytest_runtest_logstart?
# XXX count no of failures and retry N times
runner = self.config.pluginmanager.getplugin("runner") runner = self.config.pluginmanager.getplugin("runner")
fspath = nodeid.split("::")[0] fspath = nodeid.split("::")[0]
msg = "Slave %r crashed while running %r" %(slave.gateway.id, nodeid) msg = "Slave %r crashed while running %r" % (slave.gateway.id, nodeid)
rep = runner.TestReport(nodeid, (fspath, None, fspath), (), rep = runner.TestReport(nodeid, (fspath, None, fspath),
"failed", msg, "???") (), "failed", msg, "???")
rep.node = slave rep.node = slave
self.config.hook.pytest_runtest_logreport(report=rep) self.config.hook.pytest_runtest_logreport(report=rep)
class TerminalDistReporter: class TerminalDistReporter:
def __init__(self, config): def __init__(self, config):
self.config = config self.config = config

View File

@@ -11,6 +11,21 @@ import py, pytest
import sys import sys
import execnet import execnet
def pytest_addoption(parser):
group = parser.getgroup("xdist", "distributed and subprocess testing")
group._addoption('-f', '--looponfail',
action="store_true", dest="looponfail", default=False,
help="run tests in subprocess, wait for modified files "
"and re-run failing test set until all pass.")
def pytest_cmdline_main(config):
if config.getoption("looponfail"):
looponfail_main(config)
return 2 # looponfail only can get stop with ctrl-C anyway
def looponfail_main(config): def looponfail_main(config):
remotecontrol = RemoteControl(config) remotecontrol = RemoteControl(config)
rootdirs = config.getini("looponfailroots") rootdirs = config.getini("looponfailroots")
@@ -87,7 +102,7 @@ class RemoteControl(object):
result = self.runsession() result = self.runsession()
failures, reports, collection_failed = result failures, reports, collection_failed = result
if collection_failed: if collection_failed:
reports = ["Collection failed, keeping previous failure set"] pass # "Collection failed, keeping previous failure set"
else: else:
uniq_failures = [] uniq_failures = []
for failure in failures: for failure in failures:
@@ -109,7 +124,6 @@ def repr_pytest_looponfailinfo(failreports, rootdirs):
def init_slave_session(channel, args, option_dict): def init_slave_session(channel, args, option_dict):
import os, sys import os, sys
import py
outchannel = channel.gateway.newchannel() outchannel = channel.gateway.newchannel()
sys.stdout = sys.stderr = outchannel.makefile('w') sys.stdout = sys.stderr = outchannel.makefile('w')
channel.send(outchannel) channel.send(outchannel)
@@ -228,4 +242,3 @@ class StatRecorder:
changed = True changed = True
self.statcache = newstat self.statcache = newstat
return changed return changed

View File

@@ -1,18 +1,26 @@
import py import py
import pytest import pytest
def parse_numprocesses(s):
if s == 'auto':
import multiprocessing
return multiprocessing.cpu_count()
else:
return int(s)
def pytest_addoption(parser): def pytest_addoption(parser):
group = parser.getgroup("xdist", "distributed and subprocess testing") group = parser.getgroup("xdist", "distributed and subprocess testing")
group._addoption('-f', '--looponfail',
action="store_true", dest="looponfail", default=False,
help="run tests in subprocess, wait for modified files "
"and re-run failing test set until all pass.")
group._addoption('-n', dest="numprocesses", metavar="numprocesses", group._addoption('-n', dest="numprocesses", metavar="numprocesses",
action="store", type="int", action="store",
help="shortcut for '--dist=load --tx=NUM*popen'") type=parse_numprocesses,
group.addoption('--boxed', help="shortcut for '--dist=load --tx=NUM*popen', "
action="store_true", dest="boxed", default=False, "you can use 'auto' here for auto detection CPUs number on "
help="box each test run in a separate process (unix)") "host system")
group._addoption('--max-slave-restart', action="store", default=None,
help="maximum number of slaves that can be restarted "
"when crashed (set to zero to disable this feature)")
group._addoption('--dist', metavar="distmode", group._addoption('--dist', metavar="distmode",
action="store", choices=['load', 'each', 'no'], action="store", choices=['load', 'each', 'no'],
type="choice", dest="dist", default="no", type="choice", dest="dist", default="no",
@@ -45,32 +53,31 @@ def pytest_addoption(parser):
# ------------------------------------------------------------------------- # -------------------------------------------------------------------------
def pytest_addhooks(pluginmanager): def pytest_addhooks(pluginmanager):
from xdist import newhooks from xdist import newhooks
pluginmanager.addhooks(newhooks) # avoid warnings with pytest-2.8
method = getattr(pluginmanager, "add_hookspecs", None)
if method is None:
method = pluginmanager.addhooks
method(newhooks)
# ------------------------------------------------------------------------- # -------------------------------------------------------------------------
# distributed testing initialization # distributed testing initialization
# ------------------------------------------------------------------------- # -------------------------------------------------------------------------
def pytest_cmdline_main(config):
check_options(config)
if config.getvalue("looponfail"):
from xdist.looponfail import looponfail_main
looponfail_main(config)
return 2 # looponfail only can get stop with ctrl-C anyway
def pytest_configure(config, __multicall__): @pytest.mark.trylast
__multicall__.execute() def pytest_configure(config):
if config.getvalue("dist") != "no": if config.getoption("dist") != "no":
from xdist.dsession import DSession from xdist.dsession import DSession
session = DSession(config) session = DSession(config)
config.pluginmanager.register(session, "dsession") config.pluginmanager.register(session, "dsession")
tr = config.pluginmanager.getplugin("terminalreporter") tr = config.pluginmanager.getplugin("terminalreporter")
tr.showfspath = False tr.showfspath = False
def check_options(config): @pytest.mark.tryfirst
def pytest_cmdline_main(config):
if config.option.numprocesses: if config.option.numprocesses:
config.option.dist = "load" config.option.dist = "load"
config.option.tx = ['popen'] * int(config.option.numprocesses) config.option.tx = ['popen'] * config.option.numprocesses
if config.option.distload: if config.option.distload:
config.option.dist = "load" config.option.dist = "load"
val = config.getvalue val = config.getvalue
@@ -82,46 +89,3 @@ def check_options(config):
elif val("dist") != "no": elif val("dist") != "no":
if usepdb: if usepdb:
raise pytest.UsageError("--pdb incompatible with distributing tests.") raise pytest.UsageError("--pdb incompatible with distributing tests.")
def pytest_runtest_protocol(item):
if item.config.getvalue("boxed"):
reports = forked_run_report(item)
for rep in reports:
item.ihook.pytest_runtest_logreport(report=rep)
return True
def forked_run_report(item):
# for now, we run setup/teardown in the subprocess
# XXX optionally allow sharing of setup/teardown
from _pytest.runner import runtestprotocol
EXITSTATUS_TESTEXIT = 4
import marshal
from xdist.remote import serialize_report
from xdist.slavemanage import unserialize_report
def runforked():
try:
reports = runtestprotocol(item, log=False)
except KeyboardInterrupt:
py.std.os._exit(EXITSTATUS_TESTEXIT)
return marshal.dumps([serialize_report(x) for x in reports])
ff = py.process.ForkedFunc(runforked)
result = ff.waitfinish()
if result.retval is not None:
report_dumps = marshal.loads(result.retval)
return [unserialize_report("testreport", x) for x in report_dumps]
else:
if result.exitstatus == EXITSTATUS_TESTEXIT:
py.test.exit("forked test item %s raised Exit" %(item,))
return [report_process_crash(item, result)]
def report_process_crash(item, result):
path, lineno = item._getfslineno()
info = "%s:%s: running the test CRASHED with signal %d" %(
path, lineno, result.signal)
from _pytest import runner
call = runner.CallInfo(lambda: 0/0, "???")
call.excinfo = info
rep = runner.pytest_runtest_makereport(item, call)
return rep

View File

@@ -51,9 +51,12 @@ class SlaveInteractor:
elif name == "runtests_all": elif name == "runtests_all":
torun.extend(range(len(session.items))) torun.extend(range(len(session.items)))
self.log("items to run:", torun) self.log("items to run:", torun)
while torun: # only run if we have an item and a next item
while len(torun) >= 2:
self.run_tests(torun) self.run_tests(torun)
if name == "shutdown": if name == "shutdown":
if torun:
self.run_tests(torun)
break break
return True return True
@@ -125,6 +128,7 @@ def remote_initconfig(option_dict, args):
if __name__ == '__channelexec__': if __name__ == '__channelexec__':
channel = channel # noqa
# python3.2 is not concurrent import safe, so let's play it safe # python3.2 is not concurrent import safe, so let's play it safe
# https://bitbucket.org/hpk42/pytest/issue/347/pytest-xdist-and-python-32 # https://bitbucket.org/hpk42/pytest/issue/347/pytest-xdist-and-python-32
if sys.version_info[:2] == (3,2): if sys.version_info[:2] == (3,2):

View File

@@ -28,34 +28,32 @@ class NodeManager(object):
self.specs.append(spec) self.specs.append(spec)
self.roots = self._getrsyncdirs() self.roots = self._getrsyncdirs()
self.rsyncoptions = self._getrsyncoptions() self.rsyncoptions = self._getrsyncoptions()
self._rsynced_specs = py.builtin.set()
def rsync_roots(self): def rsync_roots(self, gateway):
""" make sure that all remote gateways """Rsync the set of roots to the node's gateway cwd."""
have the same set of roots in their
current directory.
"""
if self.roots: if self.roots:
# send each rsync root
for root in self.roots: for root in self.roots:
self.rsync(root, **self.rsyncoptions) self.rsync(gateway, root, **self.rsyncoptions)
def makegateways(self):
assert not list(self.group)
self.config.hook.pytest_xdist_setupnodes(config=self.config,
specs=self.specs)
for spec in self.specs:
gw = self.group.makegateway(spec)
self.config.hook.pytest_xdist_newgateway(gateway=gw)
def setup_nodes(self, putevent): def setup_nodes(self, putevent):
self.makegateways() self.config.hook.pytest_xdist_setupnodes(config=self.config,
self.rsync_roots() specs=self.specs)
self.trace("setting up nodes") self.trace("setting up nodes")
for gateway in self.group: nodes = []
node = SlaveController(self, gateway, self.config, putevent) for spec in self.specs:
gateway.node = node # to keep node alive nodes.append(self.setup_node(spec, putevent))
node.setup() return nodes
self.trace("started node %r" % node)
def setup_node(self, spec, putevent):
gw = self.group.makegateway(spec)
self.config.hook.pytest_xdist_newgateway(gateway=gw)
self.rsync_roots(gw)
node = SlaveController(self, gw, self.config, putevent)
gw.node = node # keep the node alive
node.setup()
self.trace("started node %r" % node)
return node
def teardown_nodes(self): def teardown_nodes(self):
self.group.terminate(self.EXIT_TIMEOUT) self.group.terminate(self.EXIT_TIMEOUT)
@@ -110,49 +108,43 @@ class NodeManager(object):
'verbose': self.config.option.verbose, 'verbose': self.config.option.verbose,
} }
def rsync(self, gateway, source, notify=None, verbose=False, ignores=None):
def rsync(self, source, notify=None, verbose=False, ignores=None): """Perform rsync to remote hosts for node."""
""" perform rsync to all remote hosts. # XXX This changes the calling behaviour of
""" # pytest_xdist_rsyncstart and pytest_xdist_rsyncfinish to
# be called once per rsync target.
rsync = HostRSync(source, verbose=verbose, ignores=ignores) rsync = HostRSync(source, verbose=verbose, ignores=ignores)
seen = py.builtin.set() spec = gateway.spec
gateways = [] if spec.popen and not spec.chdir:
for gateway in self.group: # XXX This assumes that sources are python-packages
spec = gateway.spec # and that adding the basedir does not hurt.
if spec.popen and not spec.chdir: gateway.remote_exec("""
# XXX this assumes that sources are python-packages import sys ; sys.path.insert(0, %r)
# and that adding the basedir does not hurt """ % os.path.dirname(str(source))).waitclose()
gateway.remote_exec(""" return
import sys ; sys.path.insert(0, %r) if (spec, source) in self._rsynced_specs:
""" % os.path.dirname(str(source))).waitclose() return
continue def finished():
if spec not in seen: if notify:
def finished(): notify("rsyncrootready", spec, source)
if notify: rsync.add_target_host(gateway, finished=finished)
notify("rsyncrootready", spec, source) self._rsynced_specs.add((spec, source))
rsync.add_target_host(gateway, finished=finished) self.config.hook.pytest_xdist_rsyncstart(
seen.add(spec) source=source,
gateways.append(gateway) gateways=[gateway],
if seen: )
self.config.hook.pytest_xdist_rsyncstart( rsync.send()
source=source, self.config.hook.pytest_xdist_rsyncfinish(
gateways=gateways, source=source,
) gateways=[gateway],
rsync.send() )
self.config.hook.pytest_xdist_rsyncfinish(
source=source,
gateways=gateways,
)
class HostRSync(execnet.RSync): class HostRSync(execnet.RSync):
""" RSyncer that filters out common files """ RSyncer that filters out common files
""" """
def __init__(self, sourcedir, *args, **kwargs): def __init__(self, sourcedir, *args, **kwargs):
self._synced = {} self._synced = {}
ignores= None self._ignores = kwargs.pop('ignores', None) or []
if 'ignores' in kwargs:
ignores = kwargs.pop('ignores')
self._ignores = ignores or []
super(HostRSync, self).__init__(sourcedir=sourcedir, **kwargs) super(HostRSync, self).__init__(sourcedir=sourcedir, **kwargs)
def filter(self, path): def filter(self, path):