Compare commits

...

46 Commits
1.9 ... 1.11

Author SHA1 Message Date
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
holger krekel
b0732f99a6 final 2014-01-29 13:10:20 +01:00
holger krekel
ab987d512c make pexpect specific environments so that the canonical ones also run on windows. Also add a little sleep to a pexpect-related tests in the hope it helps the test consistently passing (Bah). 2014-01-27 19:47:37 +01:00
holger krekel
ffa156ced5 refine distribution fractions and make it dependend on the duration
of the last test in a node's set.
2014-01-27 12:39:43 +01:00
holger krekel
e26d4486fc send multiple "to test" indices in one network message to a slave
and improve heuristics for sending chunks where the chunksize
depends on the number of remaining tests rather than fixed numbers.
This reduces the number of master -> node messages (but not the
reverse direction)
2014-01-27 11:37:40 +01:00
holger krekel
18a30fab7d fix issue419: work with collection indices instead of node ids.
This reduces network message size.
2014-01-27 11:37:33 +01:00
holger krekel
dd73d132b1 remove unneccesary item2nodes data structure for load scheduling 2014-01-27 11:37:27 +01:00
holger krekel
a22e0384a1 a failing test for pytest issue419 2014-01-27 11:37:26 +01:00
holger krekel
8eabc7d5df Merged in lukaszb/pytest-xdist/lukaszb/readmetxt-typo-1388766852691 (pull request #7)
README.txt typo
2014-01-04 21:41:34 +01:00
Lukasz Balcerzak
dbf7932514 README.txt typo 2014-01-03 16:34:15 +00:00
holger krekel
6355c3e8d0 add glob support for rsyncignores, add command line option to pass
additional rsyncignores. Thanks Anatoly Bubenkov.
2013-12-06 13:26:09 +01:00
holger krekel
c505f00320 Merged in paylogic/pytest-xdist/rsyncignore-option (pull request #5)
add glob support for rsyncignore. add command line option for rsyncignore
2013-12-06 13:25:19 +01:00
Anatoly Bubenkov
7e066be35e docs fixed 2013-12-06 10:21:32 +01:00
Anatoly Bubenkov
27f1d8049b merge with mainline 2013-12-06 10:20:03 +01:00
Anatoly Bubenkov
bcf1c85f44 docs fixed 2013-12-06 10:19:18 +01:00
Anatoly Bubenkov
a3fac52e85 more readability for return 2013-12-05 23:22:49 +01:00
Anatoly Bubenkov
152f965a70 fix ignores mixing 2013-12-05 16:35:56 +01:00
Anatoly Bubenkov
027b3f51ed change metavars 2013-12-05 15:31:52 +01:00
Anatoly Bubenkov
f4b86f0127 add glob support for rsyncignore. add command line option for rsyncignore 2013-12-05 15:19:31 +01:00
holger krekel
8146b27671 now that pytest_runtest_logstart is sent to the master again,
disable showing of filenames explicitely with the terminalreporter
2013-11-21 15:00:06 +01:00
holger krekel
a488cc5add merge 2013-11-19 12:21:31 +01:00
holger krekel
aad1aace52 merge fix pytest issue382 - produce "pytest_runtest_logstart" event again
in master. Thanks Aron Curzon.
2013-11-19 12:08:55 +01:00
curzona
82d357ab1a Uncomment logstart hook in remote.py and slavemanage.py 2013-11-09 12:10:13 -08:00
holger krekel
ca54911481 ignore directories and pyc files for changes (editors write tmp/swap
files etc., which also affects mtime of directory)
2013-10-07 10:28:08 +02:00
holger krekel
8e88261c15 ignore dot files for file changes (editors write tmp/swap files etc.) 2013-10-07 10:20:24 +02:00
holger krekel
6045195f16 Added tag 1.9 for changeset 5c5cb6d59e12 2013-10-04 14:35:34 +02:00
20 changed files with 897 additions and 338 deletions

View File

@@ -14,9 +14,17 @@ syntax:glob
*.class
*.orig
*.sublime-*
.Python
build/
dist/
include/
lib/
bin/
pytest_xdist.egg-info
issue/
3rdparty/
pytestdebug.log
.tox
.cache

View File

@@ -13,3 +13,5 @@ cd44a941c833c098e4899fe3d42a96703754d0d5 1.5
0d1c00018008433956aa7d93007bab6ea7de96e4 1.8
0d1c00018008433956aa7d93007bab6ea7de96e4 1.8
1d27987c267577899350a25ba5828d55d87083ad 1.8
5c5cb6d59e12e566fbb0217aea718dc31578bee1 1.9
4406fc2a6427fadc021ed7e43e7aa5032b1ea91f 1.10

View File

@@ -1,3 +1,39 @@
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
-------------------------
- add glob support for rsyncignores, add command line option to pass
additional rsyncignores. Thanks Anatoly Bubenkov.
- fix pytest issue382 - produce "pytest_runtest_logstart" event again
in master. Thanks Aron Curzon.
- fix pytest issue419 by sending/receiving indices into the test
collection instead of node ids (which are not neccessarily unique
for functions parametrized with duplicate values)
- send multiple "to test" indices in one network message to a slave
and improve heuristics for sending chunks where the chunksize
depends on the number of remaining tests rather than fixed numbers.
This reduces the number of master -> node messages (but not the
reverse direction)
1.9
-------------------------
@@ -9,14 +45,14 @@
- fix pytest issue41: re-run tests on all file changes, not just
randomly select ones like .py/.c.
- fix pytest issue347: slaves running on top of Python3.2
- fix pytest issue347: slaves running on top of Python3.2
will set PYTHONDONTWRITEYBTECODE to 1 to avoid import concurrency
bugs.
1.8
-------------------------
- fix pytest-issue93 - use the refined pytest-2.2.1 runtestprotocol
- fix pytest-issue93 - use the refined pytest-2.2.1 runtestprotocol
interface to perform eager teardowns for test items.
1.7

View File

@@ -8,10 +8,10 @@ test execution modes:
those for a combined test run. This allows to speed up
development or to use special resources of `remote machines`_.
* ``--boxed``: (not available on Windows) run each test in a boxed_
* ``--boxed``: (not available on Windows) run each test in a boxed_
subprocess to survive ``SEGFAULTS`` or otherwise dying processes
* ``--looponfail``: run your tests repeatedly in a subprocess. After each run
* ``--looponfail``: run your tests repeatedly in a subprocess. After each run
py.test waits until a file in your project changes and then re-runs
the previously failing tests. This is repeated until all tests pass
after which again a full run is performed.
@@ -33,7 +33,7 @@ Install the plugin with::
easy_install pytest-xdist
# or
pip install pytest-xdist
or use the package in develope/in-place mode with
@@ -80,7 +80,7 @@ Running tests in a boxed subprocess
+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
If you have tests involving C or C++ libraries you might have to deal
with tests crashing the process. For this case you max use the boxing
with tests crashing the process. For this case you may use the boxing
options::
py.test --boxed
@@ -91,7 +91,7 @@ running multiple processes to speed up the test run and use your CPU cores::
py.test -n3 --boxed
this would run 3 testing subprocesses in parallel which each
this would run 3 testing subprocesses in parallel which each
create new boxed subprocesses for each test.
@@ -122,6 +122,13 @@ py.test references tests as a fully qualified python
module path. **You will otherwise get strange errors**
during setup of the remote side.
You can specify multiple ``--rsyncignore`` glob-patterns
to be ignored when file are sent to the remote side.
There are also internal ignores: .*, *.pyc, *.pyo, *~
Those you cannot override using rsyncignore command-line or
ini-file option(s).
Sending tests to remote Socket Servers
+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++

View File

@@ -2,7 +2,7 @@ from setuptools import setup
setup(
name="pytest-xdist",
version='1.9',
version='1.11',
description='py.test xdist plugin for distributed testing and loop-on-failing modes',
long_description=open('README.txt').read(),
license='MIT',
@@ -13,7 +13,7 @@ setup(
packages = ['xdist'],
entry_points = {'pytest11': ['xdist = xdist.plugin'],},
zip_safe=False,
install_requires = ['execnet>=1.1', 'pytest>=2.3.5'],
install_requires = ['execnet>=1.1', 'pytest>=2.4.2', 'py>=1.4.22'],
classifiers=[
'Development Status :: 5 - Production/Stable',
'Intended Audience :: Developers',

View File

@@ -1,5 +1,6 @@
import py
import sys
import pytest
class TestDistribution:
def test_n1_pass(self, testdir):
@@ -193,7 +194,7 @@ class TestDistribution:
assert dest.join(subdir.basename).check(dir=1)
def test_data_exchange(self, testdir):
c1 = testdir.makeconftest("""
testdir.makeconftest("""
# This hook only called on master.
def pytest_configure_node(node):
node.slaveinput['a'] = 42
@@ -250,12 +251,13 @@ class TestDistribution:
def test_keyboard_interrupt_dist(self, testdir):
# xxx could be refined to check for return code
p = testdir.makepyfile("""
testdir.makepyfile("""
def test_sleep():
import time
time.sleep(10)
""")
child = testdir.spawn_pytest("-n1")
py.std.time.sleep(0.1)
child.expect(".*test session starts.*")
child.kill(2) # keyboard interrupt
child.expect(".*KeyboardInterrupt.*")
@@ -299,7 +301,7 @@ class TestDistEach:
class TestTerminalReporting:
def test_pass_skip_fail(self, testdir):
p = testdir.makepyfile("""
testdir.makepyfile("""
import py
def test_ok():
pass
@@ -310,9 +312,9 @@ class TestTerminalReporting:
""")
result = testdir.runpytest("-n1", "-v")
result.stdout.fnmatch_lines_random([
"*PASS*test_pass_skip_fail.py:2: *test_ok*",
"*SKIP*test_pass_skip_fail.py:4: *test_skip*",
"*FAIL*test_pass_skip_fail.py:6: *test_func*",
"*PASS*test_pass_skip_fail.py*test_ok*",
"*SKIP*test_pass_skip_fail.py*test_skip*",
"*FAIL*test_pass_skip_fail.py*test_func*",
])
result.stdout.fnmatch_lines([
"*def test_func():",
@@ -321,13 +323,13 @@ class TestTerminalReporting:
])
def test_fail_platinfo(self, testdir):
p = testdir.makepyfile("""
testdir.makepyfile("""
def test_func():
assert 0
""")
result = testdir.runpytest("-n1", "-v")
result.stdout.fnmatch_lines([
"*FAIL*test_fail_platinfo.py:1: *test_func*",
"*FAIL*test_fail_platinfo.py*test_func*",
"*0*Python*",
"*def test_func():",
"> assert 0",
@@ -350,7 +352,7 @@ def test_teardownfails_one_function(testdir):
@py.test.mark.xfail
def test_terminate_on_hangingnode(testdir):
p = testdir.makeconftest("""
def pytest_sessionfinishes(session):
def pytest_sessionfinish(session):
if session.nodeid == "my": # running on slave
import time
time.sleep(3)
@@ -362,7 +364,7 @@ def test_terminate_on_hangingnode(testdir):
])
@pytest.mark.xfail(reason="works if run outside test suite", run=False)
def test_session_hooks(testdir):
testdir.makeconftest("""
import sys
@@ -458,3 +460,83 @@ def test_issue34_pluginloading_in_subprocess(testdir):
result.stdout.fnmatch_lines([
"*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*",
])
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 failed node*",
"*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 failed node*",
"*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 failed node*",
"*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 failed node*",
"*Slave*crashed while running*",
"*2 failed*2 passed*",
])

View File

@@ -7,14 +7,11 @@ pytest_plugins = "pytester"
def pytest_addoption(parser):
parser.addoption('--gx',
action="append", dest="gspecs", default=None,
action="append", dest="gspecs",
help=("add a global test environment, XSpec-syntax. "))
def pytest_funcarg__specssh(request):
return getspecssh(request.config)
def getgspecs(config):
return [execnet.XSpec(spec)
for spec in config.getvalueorskip("gspecs")]
# configuration information for tests
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):
p1 = testdir.makepyfile("""
import os
@@ -13,12 +17,36 @@ def test_functional_boxed(testdir):
"*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:
def test_boxed_option_default(self, testdir):
tmpdir = testdir.tmpdir.ensure("subdir", dir=1)
config = testdir.parseconfig()
assert not config.option.boxed
py.test.importorskip("execnet")
pytest.importorskip("execnet")
config = testdir.parseconfig('-d', tmpdir)
assert not config.option.boxed

View File

@@ -4,8 +4,8 @@ from xdist.dsession import (
EachScheduling,
report_collection_diff,
)
from _pytest import main as outcome
import py
import pytest
import execnet
XSpec = execnet.XSpec
@@ -28,15 +28,12 @@ class MockNode:
self.sent = []
self.gateway = MockGateway()
def send_runtest(self, nodeid):
self.sent.append(nodeid)
def send_runtest_some(self, indices):
self.sent.extend(indices)
def send_runtest_all(self):
self.sent.append("ALL")
def sendlist(self, items):
self.sent.extend(items)
def shutdown(self):
self._shutdown=True
@@ -63,9 +60,9 @@ class TestEachScheduling:
assert sched.tests_finished()
assert node1.sent == ['ALL']
assert node2.sent == ['ALL']
sched.remove_item(node1, collection[0])
sched.remove_item(node1, 0)
assert sched.tests_finished()
sched.remove_item(node2, collection[0])
sched.remove_item(node2, 0)
assert sched.tests_finished()
def test_schedule_remove_node(self):
@@ -86,11 +83,10 @@ class TestEachScheduling:
class TestLoadScheduling:
def test_schedule_load_simple(self):
node1 = MockNode()
node2 = MockNode()
sched = LoadScheduling(2)
sched.addnode(node1)
sched.addnode(node2)
sched.addnode(MockNode())
sched.addnode(MockNode())
node1, node2 = sched.nodes
collection = ["a.py::test_1", "a.py::test_2"]
assert not sched.collection_is_completed
sched.addnode_collection(node1, collection)
@@ -100,39 +96,40 @@ class TestLoadScheduling:
assert sched.node2collection[node1] == collection
assert sched.node2collection[node2] == collection
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 == collection
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.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):
sched = LoadScheduling(2)
node1 = MockNode()
node2 = MockNode()
sched.addnode(node1)
sched.addnode(node2)
sched.ITEM_CHUNKSIZE = 2
col = ["xyz"] * (2*sched.ITEM_CHUNKSIZE +1)
sched.addnode(MockNode())
sched.addnode(MockNode())
node1, node2 = sched.nodes
col = ["xyz"] * (6)
sched.addnode_collection(node1, col)
sched.addnode_collection(node2, col)
sched.init_distribute()
#assert not sched.tests_finished()
sent1 = node1.sent
sent2 = node2.sent
chunkitems = col[:sched.ITEM_CHUNKSIZE]
assert sent1 == chunkitems
assert sent2 == chunkitems
assert sent1 == [0, 1]
assert sent2 == [2, 3]
assert sched.pending == [4, 5]
assert sched.node2pending[node1] == sent1
assert sched.node2pending[node2] == sent2
assert len(sched.pending) == 1
for node in (node1, node2):
for i in range(sched.ITEM_CHUNKSIZE):
sched.remove_item(node, "xyz")
assert len(sched.pending) == 2
sched.remove_item(node1, 0)
assert node1.sent == [0, 1, 4]
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
def test_add_remove_node(self):
@@ -147,6 +144,25 @@ class TestLoadScheduling:
crashitem = sched.remove_node(node)
assert crashitem == collection[0]
def test_schedule_different_tests_collected(self):
"""
Test that LoadScheduling is logging different tests were
collected by slaves when that happens.
"""
node1 = MockNode()
node2 = MockNode()
sched = LoadScheduling(2)
logged_messages = []
py.log.setconsumer('loadsched', logged_messages.append)
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()
logged_content = ''.join(x.content() for x in logged_messages)
assert 'Different tests were collected between' in logged_content
assert 'Different tests collected, aborting run' in logged_content
class TestDistReporter:
@@ -182,7 +198,7 @@ class TestDistReporter:
def test_report_collection_diff_equal():
"""Test reporting of equal collections."""
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():
@@ -205,7 +221,18 @@ def test_report_collection_diff_different():
'-YYY'
)
try:
report_collection_diff(from_collection, to_collection, 1, 2)
except AssertionError as e:
assert py.builtin._totext(e) == error_message
msg = report_collection_diff(from_collection, to_collection, 1, 2)
assert msg == error_message
@pytest.mark.xfail(reason="duplicate test ids not supported yet")
def test_pytest_issue419(testdir):
testdir.makepyfile("""
import pytest
@pytest.mark.parametrize('birth_year', [1988, 1988, ])
def test_2011_table(birth_year):
pass
""")
reprec = testdir.inline_run("-n1")
reprec.assertoutcome(passed=2)
assert 0

View File

@@ -14,6 +14,10 @@ class TestStatRecorder:
changed = sd.check()
assert changed
(hello + "c").write("hello")
changed = sd.check()
assert not changed
p = tmp.ensure("new.py")
changed = sd.check()
assert changed
@@ -36,6 +40,12 @@ class TestStatRecorder:
changed = sd.check()
assert changed
def test_dirchange(self, tmpdir):
tmp = tmpdir
tmp.ensure("dir", "hello.py")
sd = StatRecorder([tmp])
assert not sd.fil(tmp.join("dir"))
def test_filechange_deletion_race(self, tmpdir, monkeypatch):
tmp = tmpdir
sd = StatRecorder([tmp])
@@ -63,12 +73,10 @@ class TestStatRecorder:
pycfile = hello + "c"
pycfile.ensure()
changed = sd.check()
assert changed
hello.write("world")
changed = sd.check()
assert changed
assert not pycfile.check()
def test_waitonchange(self, tmpdir, monkeypatch):
tmp = tmpdir

View File

@@ -46,6 +46,11 @@ class TestDistOptions:
assert nm.roots
assert testdir.tmpdir in nm.roots
def test_getrsyncignore(self, testdir):
config = testdir.parseconfigure('--rsyncignore=fo*')
nm = NodeManager(config, specs=[execnet.XSpec("popen//chdir=qwe")])
assert 'fo*' in nm.rsyncoptions['ignores']
def test_getrsyncdirs_with_conftest(self, testdir):
p = py.path.local()
for bn in 'x y z'.split():

View File

@@ -3,7 +3,6 @@ from xdist.slavemanage import SlaveController, unserialize_report
from xdist.remote import serialize_report
import execnet
queue = py.builtin._tryimport("queue", "Queue")
from py.builtin import print_
import marshal
WAIT_TIMEOUT = 10.0
@@ -26,7 +25,7 @@ class SlaveSetup:
use_callback = False
def __init__(self, request):
self.testdir = testdir = request.getfuncargvalue("testdir")
self.testdir = request.getfuncargvalue("testdir")
self.request = request
self.events = queue.Queue()
@@ -140,7 +139,7 @@ class TestReportSerialization:
class TestSlaveInteractor:
def test_basic_collect_and_runtests(self, slave):
p = slave.testdir.makepyfile("""
slave.testdir.makepyfile("""
def test_func():
pass
""")
@@ -154,8 +153,11 @@ class TestSlaveInteractor:
assert ev.kwargs['topdir'] == slave.testdir.tmpdir
ids = ev.kwargs['ids']
assert len(ids) == 1
slave.sendcommand("runtests", ids=ids)
slave.sendcommand("runtests", indices=list(range(len(ids))))
slave.sendcommand("shutdown")
ev = slave.popevent("logstart")
assert ev.kwargs["nodeid"].endswith("test_func")
assert len(ev.kwargs["location"]) == 3
ev = slave.popevent("testreport") # setup
ev = slave.popevent("testreport")
assert ev.name == "testreport"
@@ -167,7 +169,7 @@ class TestSlaveInteractor:
assert 'slaveoutput' in ev.kwargs
def test_remote_collect_skip(self, slave):
p = slave.testdir.makepyfile("""
slave.testdir.makepyfile("""
import py
py.test.skip("hello")
""")
@@ -184,7 +186,7 @@ class TestSlaveInteractor:
assert not ev.kwargs['ids']
def test_remote_collect_fail(self, slave):
p = slave.testdir.makepyfile("""aasd qwe""")
slave.testdir.makepyfile("""aasd qwe""")
slave.setup()
ev = slave.popevent("collectionstart")
assert not ev.kwargs
@@ -198,7 +200,7 @@ class TestSlaveInteractor:
assert not ev.kwargs['ids']
def test_runtests_all(self, slave):
p = slave.testdir.makepyfile("""
slave.testdir.makepyfile("""
def test_func(): pass
def test_func2(): pass
""")
@@ -225,7 +227,7 @@ class TestSlaveInteractor:
def test_happy_run_events_converted(self, testdir, slave):
py.test.xfail("implement a simple test for event production")
assert not slave.use_callback
p = slave.testdir.makepyfile("""
slave.testdir.makepyfile("""
def test_func():
pass
""")

View File

@@ -1,6 +1,7 @@
import py
import os
import pytest
import execnet
from xdist import slavemanage
from xdist.slavemanage import HostRSync, NodeManager
pytest_plugins = "pytester",
@@ -24,6 +25,14 @@ def pytest_funcarg__mysetup(request):
request.getfuncargvalue("_pytest")
return mysetup(request)
@pytest.fixture
def slavecontroller(monkeypatch):
class MockController(object):
def __init__(self, *args): pass
def setup(self): pass
monkeypatch.setattr(slavemanage, 'SlaveController', MockController)
return MockController
class TestNodeManagerPopen:
def test_popen_no_default_chdir(self, config):
gm = NodeManager(config, ["popen"])
@@ -36,9 +45,10 @@ class TestNodeManagerPopen:
for spec in NodeManager(config, l, defaultchdir="abc").specs:
assert spec.chdir == "abc"
def test_popen_makegateway_events(self, config, hookrecorder, _pytest):
def test_popen_makegateway_events(self, config,
hookrecorder, _pytest, slavecontroller):
hm = NodeManager(config, ["popen"] * 2)
hm.makegateways()
hm.setup_nodes(None)
call = hookrecorder.popcall("pytest_xdist_setupnodes")
assert len(call.specs) == 2
@@ -51,10 +61,10 @@ class TestNodeManagerPopen:
hm.teardown_nodes()
assert not len(hm.group)
def test_popens_rsync(self, config, mysetup):
def test_popens_rsync(self, config, mysetup, slavecontroller):
source = mysetup.source
hm = NodeManager(config, ["popen"] * 2)
hm.makegateways()
hm.setup_nodes(None)
assert len(hm.group) == 2
for gw in hm.group:
class pseudoexec:
@@ -65,19 +75,21 @@ class TestNodeManagerPopen:
pass
gw.remote_exec = pseudoexec
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
hm.teardown_nodes()
assert not len(hm.group)
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
hm = NodeManager(config, ["popen//chdir=%s" %dest] * 1)
hm.makegateways()
hm = NodeManager(config, ["popen//chdir=%s" % dest] * 1)
hm.setup_nodes(None)
source.ensure("dir1", "dir2", "hello")
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 l[0] == ("rsyncrootready", hm.group['gw0'].spec, source)
hm.teardown_nodes()
@@ -86,12 +98,15 @@ class TestNodeManagerPopen:
assert dest.join("dir1", "dir2").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
hm = NodeManager(config, ["popen//chdir=%s" %dest] * 2)
hm.makegateways()
hm = NodeManager(config, ["popen//chdir=%s" % dest] * 2)
hm.roots = []
hm.setup_nodes(None)
source.ensure("dir1", "dir2", "hello")
hm.rsync(source)
gw = hm.group[0]
hm.rsync(gw, source)
call = hookrecorder.popcall("pytest_xdist_rsyncstart")
assert call.source == source
assert len(call.gateways) == 1
@@ -108,12 +123,12 @@ class TestHRSync:
return mysetup(request)
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(".svn", "entries")
source.ensure(".somedotfile", "moreentries")
source.ensure("somedir", "editfile~")
syncer = HostRSync(source)
syncer = HostRSync(source, ignores=NodeManager.DEFAULT_IGNORES)
l = list(source.visit(rec=syncer.filter,
fil=syncer.filter))
assert len(l) == 3
@@ -139,7 +154,7 @@ class TestNodeManager:
@py.test.mark.xfail
def test_rsync_roots_no_roots(self, testdir, mysetup):
mysetup.source.ensure("dir1", "file1").write("hello")
config = testdir.parseconfig(source)
config = testdir.parseconfig(mysetup.source)
nodemanager = NodeManager(config, ["popen//chdir=%s" % mysetup.dest])
#assert nodemanager.config.topdir == source == config.topdir
nodemanager.makegateways()
@@ -152,7 +167,7 @@ class TestNodeManager:
assert p.join("dir1").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
dir1 = mysetup.source.mkdir("dir1")
dir2 = dir1.mkdir("dir2")
@@ -164,8 +179,7 @@ class TestNodeManager:
"--rsyncdir", rsyncroot,
source,
))
nodemanager.makegateways()
nodemanager.rsync_roots()
nodemanager.setup_nodes(None) # calls .rsync_roots()
if rsyncroot == source:
dest = dest.join("source")
assert dest.join("dir1").check()
@@ -173,7 +187,7 @@ class TestNodeManager:
assert dest.join("dir1", "dir2", 'hello').check()
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
dir2 = source.ensure("dir1", "dir2", dir=1)
source.ensure("dir1", "somefile", dir=1)
@@ -185,41 +199,43 @@ class TestNodeManager:
"""))
config = testdir.parseconfig(source)
nodemanager = NodeManager(config, ["popen//chdir=%s" % dest])
nodemanager.makegateways()
nodemanager.rsync_roots()
nodemanager.setup_nodes(None) # calls .rsync_roots()
assert dest.join("dir2").check()
assert not dest.join("dir1").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
dir2 = source.ensure("dir1", "dir2", dir=1)
dir5 = source.ensure("dir5", "dir6", "bogus")
dirf = source.ensure("dir5", "file")
source.ensure("dir5", "dir6", "bogus")
source.ensure("dir5", "file")
dir2.ensure("hello")
source.ensure("foo", "bar")
source.ensure("bar", "foo")
source.join("tox.ini").write(py.std.textwrap.dedent("""
[pytest]
rsyncdirs = dir1 dir5
rsyncignore = dir1/dir2 dir5/dir6
rsyncignore = dir1/dir2 dir5/dir6 foo*
"""))
config = testdir.parseconfig(source)
config.option.rsyncignore = ['bar']
nodemanager = NodeManager(config, ["popen//chdir=%s" % dest])
nodemanager.makegateways()
nodemanager.rsync_roots()
nodemanager.setup_nodes(None) # calls .rsync_roots()
assert dest.join("dir1").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('foo').check()
assert not dest.join('bar').check()
def test_optimise_popen(self, testdir, mysetup):
source, dest = mysetup.source, mysetup.dest
def test_optimise_popen(self, testdir, mysetup, slavecontroller):
source = mysetup.source
specs = ["popen"] * 3
source.join("conftest.py").write("rsyncdirs = ['a']")
source.ensure('a', dir=1)
config = testdir.parseconfig(source)
nodemanager = NodeManager(config, specs)
nodemanager.makegateways()
nodemanager.rsync_roots()
nodemanager.setup_nodes(None) # calls .rysnc_roots()
for gwspec in nodemanager.specs:
assert gwspec._samefilesystem()
assert not gwspec.chdir
@@ -230,8 +246,6 @@ class TestNodeManager:
pass
""")
reprec = testdir.inline_run("-d", "--rsyncdir=%s" % testdir.tmpdir,
"--tx", specssh, testdir.tmpdir)
"--tx", specssh, testdir.tmpdir)
rep, = reprec.getreports("pytest_runtest_logreport")
assert rep.passed

35
tox.ini
View File

@@ -1,25 +1,42 @@
[tox]
envlist=py26,py32,py33,py27,py26,py26-old,py33-old
envlist=py26,py33,py34,py27,py27-pexpect,py33-pexpect,py26,py26-old,py33-old,flakes
[testenv]
changedir=testing
deps=pytest>=2.4.2
commands= py.test --junitxml={envlogdir}/junit-{envname}.xml []
deps=pytest>=2.5.1
commands= py.test {posargs}
[testenv:py27]
deps=
pytest>=2.4.2
[testenv:py27-pexpect]
deps={[testenv]deps}
pexpect
[testenv:py33-pexpect]
deps={[testenv]deps}
pexpect
[testenv:flakes]
changedir=
deps = pytest-flakes>=0.2
commands = py.test --flakes -m flakes testing xdist
[testenv:py26-old]
basepython = python2.6
deps=
pytest==2.3.5
pexpect
pytest==2.5.2
pycmd
commands=
py.cleanup -a
py.test {posargs}
[testenv:py33-old]
basepython = python3.3
deps=
pytest==2.3.5
pytest==2.5.2
pycmd
commands=
py.cleanup -a
py.test {posargs}
[pytest]
addopts = -rsfxX

View File

@@ -1,2 +1,2 @@
#
__version__ = '1.9'
__version__ = '1.11'

View File

@@ -1,4 +1,3 @@
import sys
import difflib
import pytest
@@ -10,166 +9,401 @@ queue = py.builtin._tryimport('queue', 'Queue')
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):
self.numnodes = numnodes
self.node2collection = {}
self.node2pending = {}
self._started = []
self._removed2pending = {}
if log is None:
self.log = py.log.Producer("eachsched")
else:
self.log = log.loadsched
self.log = log.eachsched
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):
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):
self.node2collection[node] = None
assert node not in self.node2pending
self.node2pending[node] = []
def tests_finished(self):
if not self.collection_is_completed:
return False
if self._removed2pending:
return False
for pending in self.node2pending.values():
if len(pending) >= 2:
return False
return True
def addnode_collection(self, node, collection):
assert not self.collection_is_completed
assert self.node2collection[node] is None
self.node2collection[node] = list(collection)
self.node2pending[node] = []
if len(self.node2pending) >= self.numnodes:
self.collection_is_completed = True
"""Add the collected test items from a node
def remove_item(self, node, item):
self.node2pending[node].remove(item)
Collection is complete once all nodes have submitted their
collection. In this case it's peding list is set to an empty
list. When the collection is already completed this
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:
if collection != self.node2collection[deadnode]:
msg = report_collection_diff(self.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):
self.node2pending[node].remove(item_index)
def remove_node(self, node):
# KeyError if we didn't get an addnode() yet
pending = self.node2pending.pop(node)
if not pending:
return
crashitem = pending.pop(0)
# XXX what about the rest of pending?
crashitem = self.node2collection[node][pending.pop(0)]
if pending:
self._removed2pending[node] = pending
return crashitem
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
for node, pending in self.node2pending.items():
node.send_runtest_all()
pending[:] = self.node2collection[node]
if node in self._started:
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:
LOAD_THRESHOLD_NEWITEMS = 5
ITEM_CHUNKSIZE = 10
"""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.
"""
def __init__(self, numnodes, log=None):
self.numnodes = numnodes
self.node2pending = {}
self.node2collection = {}
self.node2pending = {}
self.pending = []
self.collection = None
if log is None:
self.log = py.log.Producer("loadsched")
else:
self.log = log.loadsched
self.collection_is_completed = False
@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):
"""Return True if nodes exist in the scheduler."""
return bool(self.node2pending)
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] = []
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
#for items in self.node2pending.values():
# if items:
# return False
if self.pending:
return False
for pending in self.node2pending.values():
if len(pending) >= 2:
return False
return True
def addnode_collection(self, node, collection):
assert not self.collection_is_completed
assert node in self.node2pending
self.node2collection[node] = list(collection)
if len(self.node2collection) >= self.numnodes:
self.collection_is_completed = True
"""Add the collected test items from a node
def remove_item(self, node, item):
if item not in self.item2nodes:
raise AssertionError(item, self.item2nodes)
nodes = self.item2nodes[item]
if node in nodes: # the node might have gone down already
nodes.remove(node)
#if not nodes:
# del self.item2nodes[item]
pending = self.node2pending[node]
pending.remove(item)
# pre-load items-to-test if the node may become ready
if self.pending and len(pending) < self.LOAD_THRESHOLD_NEWITEMS:
item = self.pending.pop(0)
pending.append(item)
self.item2nodes.setdefault(item, []).append(node)
node.send_runtest(item)
self.log("items waiting for node: %d" %(len(self.pending)))
#self.log("item2pending still executing: %s" %(self.item2nodes,))
#self.log("node2pending: %s" %(self.node2pending,))
The collection is stored in the ``.node2collection`` map.
Called by the ``DSession.slave_collectionfinish`` hook.
"""
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)
def remove_item(self, node, item_index, duration=0):
"""Mark test item as completed by node
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:
# how many nodes do we have?
num_nodes = len(self.node2pending)
# if our node goes below a heuristic minimum, fill it out to
# heuristic maximum
items_per_node_min = max(2, len(self.pending) // num_nodes // 4)
items_per_node_max = max(2, len(self.pending) // num_nodes // 2)
node_pending = self.node2pending[node]
if len(node_pending) < items_per_node_min:
if duration >= 0.1 and len(node_pending) >= 2:
# 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.log("num items waiting for node:", len(self.pending))
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)
# KeyError if we didn't get an addnode() yet
for item in pending:
l = self.item2nodes[item]
l.remove(node)
if not l:
del self.item2nodes[item]
if not pending:
return
crashitem = pending.pop(0)
# The node crashed, reassing pending items
crashitem = self.collection[pending.pop(0)]
self.pending.extend(pending)
for node in self.node2pending:
self.check_schedule(node)
return crashitem
def init_distribute(self):
"""Initiate distribution of the test collection
Initiate scheduling of the items across the nodes. If this
gets called again later it behaves the same as calling
``.check_schedule()`` on all nodes so that newly added nodes
will start to be used.
This is called by the ``DSession.slave_collectionfinish`` hook
if ``.collection_is_completed`` is True.
XXX Perhaps this method should have been called ".schedule()".
"""
assert self.collection_is_completed
assert not hasattr(self, 'item2nodes')
self.item2nodes = {}
# Initial distribution already happend, reschedule on all nodes
if self.collection is not None:
for node in self.nodes:
self.check_schedule(node)
return
# XXX allow nodes to have different collections
first_node, col = list(self.node2collection.items())[0]
for node, collection in self.node2collection.items():
report_collection_diff(
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?
items_per_node = len(self.collection) // len(self.node2pending)
# take a fraction of tests for initial distribution
node_chunksize = max(items_per_node // 4, 2)
# and initialize each node with a chunk of tests
for node in self.nodes:
self._send_tests(node, node_chunksize)
def _send_tests(self, node, num):
tests_per_node = self.pending[:num]
if tests_per_node:
del self.pending[:num]
self.node2pending[node].extend(tests_per_node)
node.send_runtest_some(tests_per_node)
def _check_nodes_have_same_collection(self):
"""Return True if all nodes have collected the same items.
If collections differ this returns False and logs the
collection differences as they are found.
"""
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:
self.log(msg)
same_collection = False
self.pending = col
if not col:
return
available = list(self.node2pending.items())
num_available = self.numnodes
max_one_round = num_available * self.ITEM_CHUNKSIZE - 1
for i, item in enumerate(self.pending):
nodeindex = i % num_available
node, pending = available[nodeindex]
node.send_runtest(item)
self.item2nodes.setdefault(item, []).append(node)
pending.append(item)
if i >= max_one_round:
break
del self.pending[:i + 1]
return same_collection
def report_collection_diff(from_collection, to_collection, from_id, to_id):
"""Report the collected test difference between two nodes.
:returns: True if collections are equal.
:raises: AssertionError with a detailed error message describing the
difference between the collections.
:returns: detailed message describing the difference between the given
collections, or None if they are equal.
"""
if from_collection == to_collection:
return True
return None
diff = difflib.unified_diff(
from_collection,
@@ -183,13 +417,26 @@ def report_collection_diff(from_collection, to_collection, from_id, to_id):
'{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")])
raise AssertionError(msg)
return msg
class Interrupted(KeyboardInterrupt):
""" signals an immediate interruption. """
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):
self.config = config
self.log = py.log.Producer("dsession")
@@ -200,6 +447,7 @@ class DSession:
self.maxfail = config.getvalue("maxfail")
self.queue = queue.Queue()
self._failed_collection_errors = {}
self._active_nodes = set()
try:
self.terminal = config.pluginmanager.getplugin("terminalreporter")
except KeyError:
@@ -208,18 +456,33 @@ class DSession:
self.trdist = TerminalDistReporter(config)
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):
if self.terminal and self.config.option.verbose >= 0:
self.terminal.write_line(line)
@pytest.mark.trylast
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.setup_nodes(putevent=self.queue.put)
nodes = self.nodemanager.setup_nodes(putevent=self.queue.put)
self._active_nodes.update(nodes)
def pytest_sessionfinish(self, session):
""" teardown any resources after a test run. """
nm = getattr(self, 'nodemanager', None) # if not fully initialized
"""Shutdown all nodes."""
nm = getattr(self, 'nodemanager', None) # if not fully initialized
if nm is not None:
nm.teardown_nodes()
@@ -237,7 +500,6 @@ class DSession:
else:
assert 0, dist
self.shouldstop = False
self.session_finished = False
while not self.session_finished:
self.loop_once()
if self.shouldstop:
@@ -245,7 +507,7 @@ class DSession:
return True
def loop_once(self):
""" process one callback from one of the slaves. """
"""Process one callback from one of the slaves."""
while 1:
try:
eventcall = self.queue.get(timeout=2.0)
@@ -256,7 +518,7 @@ class DSession:
assert callname, kwargs
method = "slave_" + callname
call = getattr(self, method)
self.log("calling method: %s(**%s)" % (method, kwargs))
self.log("calling method", method, kwargs)
call(**kwargs)
if self.sched.tests_finished():
self.triggershutdown()
@@ -266,26 +528,40 @@ class DSession:
#
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['id'] = node.gateway.id
node.slaveinfo['spec'] = node.gateway.spec
self.config.hook.pytest_testnodeready(node=node)
self.sched.addnode(node)
if self.shuttingdown:
node.shutdown()
else:
self.sched.addnode(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)
if node.slaveoutput['exitstatus'] == 2: # keyboard-interrupt
if node.slaveoutput['exitstatus'] == 2: # keyboard-interrupt
self.shouldstop = "%s received keyboard-interrupt" % (node,)
self.slave_errordown(node, "keyboard-interrupt")
return
crashitem = self.sched.remove_node(node)
#assert not crashitem, (crashitem, node)
if self.shuttingdown and not self.sched.hasnodes():
self.session_finished = True
if node in self.sched.nodes:
crashitem = self.sched.remove_node(node)
assert not crashitem, (crashitem, node)
self._active_nodes.remove(node)
def slave_errordown(self, node, error):
"""Emitted by the SlaveController when a node dies."""
self.config.hook.pytest_testnodedown(node=node, error=error)
try:
crashitem = self.sched.remove_node(node)
@@ -294,41 +570,70 @@ class DSession:
else:
if crashitem:
self.handle_crashitem(crashitem, node)
#self.report_line("item crashed on node: %s" % crashitem)
if not self.sched.hasnodes():
self.session_finished = True
self.report_line("Replacing failed node %s" % node.gateway.id)
self._clone_node(node)
self._active_nodes.remove(node)
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
self.sched.addnode_collection(node, ids)
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.terminal:
if self.terminal and not self.sched.haspending():
self.trdist.ensure_show_status()
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.init_distribute()
def slave_logstart(self, node, nodeid, location):
"""Emitted when a node calls the pytest_runtest_logstart hook."""
self.config.hook.pytest_runtest_logstart(
nodeid=nodeid, location=location)
def slave_testreport(self, node, rep):
if not (rep.passed and rep.when != "call"):
if rep.when in ("setup", "call"):
self.sched.remove_item(node, rep.nodeid)
"""Emitted when a node calls the pytest_runtest_logreport hook.
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))
rep.node = node
self.config.hook.pytest_runtest_logreport(report=rep)
self._handlefailures(rep)
def slave_collectreport(self, node, rep):
"""Emitted when a node calls the pytest_collectreport hook."""
if rep.failed:
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):
# Check we haven't already seen this report (from
# another slave).
@@ -347,19 +652,21 @@ class DSession:
def triggershutdown(self):
self.log("triggering shutdown")
self.shuttingdown = True
for node in self.sched.node2pending:
for node in self.sched.nodes:
node.shutdown()
def handle_crashitem(self, nodeid, slave):
# 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")
fspath = nodeid.split("::")[0]
msg = "Slave %r crashed while running %r" %(slave.gateway.id, nodeid)
rep = runner.TestReport(nodeid, (fspath, None, fspath), (),
"failed", msg, "???")
msg = "Slave %r crashed while running %r" % (slave.gateway.id, nodeid)
rep = runner.TestReport(nodeid, (fspath, None, fspath),
(), "failed", msg, "???")
rep.node = slave
self.config.hook.pytest_runtest_logreport(report=rep)
class TerminalDistReporter:
def __init__(self, config):
self.config = config

View File

@@ -87,7 +87,7 @@ class RemoteControl(object):
result = self.runsession()
failures, reports, collection_failed = result
if collection_failed:
reports = ["Collection failed, keeping previous failure set"]
pass # "Collection failed, keeping previous failure set"
else:
uniq_failures = []
for failure in failures:
@@ -109,7 +109,6 @@ def repr_pytest_looponfailinfo(failreports, rootdirs):
def init_slave_session(channel, args, option_dict):
import os, sys
import py
outchannel = channel.gateway.newchannel()
sys.stdout = sys.stderr = outchannel.makefile('w')
channel.send(outchannel)
@@ -188,7 +187,7 @@ class StatRecorder:
self.check() # snapshot state
def fil(self, p):
return True # we are sensitive to all file changes since 1.9
return p.check(file=1, dotfile=0) and p.ext != ".pyc"
def rec(self, p):
return p.check(dotfile=0)

View File

@@ -1,5 +1,5 @@
import sys
import py, pytest
import py
import pytest
def pytest_addoption(parser):
group = parser.getgroup("xdist", "distributed and subprocess testing")
@@ -28,12 +28,14 @@ def pytest_addoption(parser):
group._addoption('-d',
action="store_true", dest="distload", default=False,
help="load-balance tests. shortcut for '--dist=load'")
group.addoption('--rsyncdir', action="append", default=[], metavar="dir1",
group.addoption('--rsyncdir', action="append", default=[], metavar="DIR",
help="add directory for rsyncing to remote tx nodes.")
group.addoption('--rsyncignore', action="append", default=[], metavar="GLOB",
help="add expression for ignores when rsyncing to remote tx nodes.")
parser.addini('rsyncdirs', 'list of (relative) paths to be rsynced for'
' remote distributed testing.', type="pathlist")
parser.addini('rsyncignore', 'list of (relative) paths to be ignored '
parser.addini('rsyncignore', 'list of (relative) glob-style paths to be ignored '
'for rsyncing.', type="pathlist")
parser.addini("looponfailroots", type="pathlist",
help="directories to check for changes", default=[py.path.local()])
@@ -51,17 +53,19 @@ def pytest_addhooks(pluginmanager):
def pytest_cmdline_main(config):
check_options(config)
if config.getvalue("looponfail"):
if config.getoption("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__):
__multicall__.execute()
if config.getvalue("dist") != "no":
if config.getoption("dist") != "no":
from xdist.dsession import DSession
session = DSession(config)
config.pluginmanager.register(session, "dsession")
tr = config.pluginmanager.getplugin("terminalreporter")
tr.showfspath = False
def check_options(config):
if config.option.numprocesses:
@@ -114,10 +118,14 @@ def forked_run_report(item):
def report_process_crash(item, result):
path, lineno = item._getfslineno()
info = "%s:%s: running the test CRASHED with signal %d" %(
path, lineno, result.signal)
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

@@ -24,7 +24,7 @@ class SlaveInteractor:
def pytest_internalerror(self, excrepr):
for line in str(excrepr).split("\n"):
self.log("IERROR> " + line)
self.log("IERROR>", line)
def pytest_sessionstart(self, session):
self.session = session
@@ -45,41 +45,44 @@ class SlaveInteractor:
torun = []
while 1:
name, kwargs = self.channel.receive()
self.log("received command %s(**%s)" % (name, kwargs))
self.log("received command", name, kwargs)
if name == "runtests":
ids = kwargs['ids']
for nodeid in ids:
torun.append(self._id2item[nodeid])
torun.extend(kwargs['indices'])
elif name == "runtests_all":
torun.extend(session.items)
self.log("items to run: %s" %(len(torun)))
torun.extend(range(len(session.items)))
self.log("items to run:", torun)
# only run if we have an item and a next item
while len(torun) >= 2:
item = torun.pop(0)
nextitem = torun[0]
self.config.hook.pytest_runtest_protocol(item=item,
nextitem=nextitem)
self.run_tests(torun)
if name == "shutdown":
while torun:
self.config.hook.pytest_runtest_protocol(
item=torun.pop(0), nextitem=None)
if torun:
self.run_tests(torun)
break
return True
def run_tests(self, torun):
items = self.session.items
self.item_index = torun.pop(0)
if torun:
nextitem = items[torun[0]]
else:
nextitem = None
self.config.hook.pytest_runtest_protocol(
item=items[self.item_index],
nextitem=nextitem)
def pytest_collection_finish(self, session):
self._id2item = {}
ids = []
for item in session.items:
self._id2item[item.nodeid] = item
ids.append(item.nodeid)
self.sendevent("collectionfinish",
topdir=str(session.fspath),
ids=ids)
ids=[item.nodeid for item in session.items])
#def pytest_runtest_logstart(self, nodeid, location, fspath):
# self.sendevent("logstart", nodeid=nodeid, location=location)
def pytest_runtest_logstart(self, nodeid, location):
self.sendevent("logstart", nodeid=nodeid, location=location)
def pytest_runtest_logreport(self, report):
data = serialize_report(report)
data["item_index"] = self.item_index
assert self.session.items[self.item_index].nodeid == report.nodeid
self.sendevent("testreport", data=data)
def pytest_collectreport(self, report):
@@ -125,6 +128,7 @@ def remote_initconfig(option_dict, args):
if __name__ == '__channelexec__':
channel = channel # noqa
# 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
if sys.version_info[:2] == (3,2):

View File

@@ -1,5 +1,8 @@
import py, pytest
import sys, os
import fnmatch
import os
import py
import pytest
import execnet
import xdist.remote
@@ -7,6 +10,7 @@ from _pytest import runner # XXX load dynamically
class NodeManager(object):
EXIT_TIMEOUT = 10
DEFAULT_IGNORES = ['.*', '*.pyc', '*.pyo', '*~']
def __init__(self, config, specs=None, defaultchdir="pyexecnetcache"):
self.config = config
self._nodesready = py.std.threading.Event()
@@ -23,38 +27,33 @@ class NodeManager(object):
self.group.allocate_id(spec)
self.specs.append(spec)
self.roots = self._getrsyncdirs()
self.rsyncoptions = self._getrsyncoptions()
self._rsynced_specs = py.builtin.set()
def rsync_roots(self):
""" make sure that all remote gateways
have the same set of roots in their
current directory.
"""
options = {
'ignores': self.config.getini("rsyncignore"),
'verbose': self.config.option.verbose,
}
def rsync_roots(self, gateway):
"""Rsync the set of roots to the node's gateway cwd."""
if self.roots:
# send each rsync root
for root in self.roots:
self.rsync(root, **options)
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)
self.rsync(gateway, root, **self.rsyncoptions)
def setup_nodes(self, putevent):
self.makegateways()
self.rsync_roots()
self.config.hook.pytest_xdist_setupnodes(config=self.config,
specs=self.specs)
self.trace("setting up nodes")
for gateway in self.group:
node = SlaveController(self, gateway, self.config, putevent)
gateway.node = node # to keep node alive
node.setup()
self.trace("started node %r" % node)
nodes = []
for spec in self.specs:
nodes.append(self.setup_node(spec, putevent))
return nodes
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):
self.group.terminate(self.EXIT_TIMEOUT)
@@ -98,38 +97,47 @@ class NodeManager(object):
roots.append(root)
return roots
def rsync(self, source, notify=None, verbose=False, ignores=None):
""" perform rsync to all remote hosts.
"""
def _getrsyncoptions(self):
"""Get options to be passed for rsync."""
ignores = list(self.DEFAULT_IGNORES)
ignores += self.config.option.rsyncignore
ignores += self.config.getini("rsyncignore")
return {
'ignores': ignores,
'verbose': self.config.option.verbose,
}
def rsync(self, gateway, source, notify=None, verbose=False, ignores=None):
"""Perform rsync to remote hosts for node."""
# 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)
seen = py.builtin.set()
gateways = []
for gateway in self.group:
spec = gateway.spec
if spec.popen and not spec.chdir:
# XXX this assumes that sources are python-packages
# and that adding the basedir does not hurt
gateway.remote_exec("""
import sys ; sys.path.insert(0, %r)
""" % os.path.dirname(str(source))).waitclose()
continue
if spec not in seen:
def finished():
if notify:
notify("rsyncrootready", spec, source)
rsync.add_target_host(gateway, finished=finished)
seen.add(spec)
gateways.append(gateway)
if seen:
self.config.hook.pytest_xdist_rsyncstart(
source=source,
gateways=gateways,
)
rsync.send()
self.config.hook.pytest_xdist_rsyncfinish(
source=source,
gateways=gateways,
)
spec = gateway.spec
if spec.popen and not spec.chdir:
# XXX This assumes that sources are python-packages
# and that adding the basedir does not hurt.
gateway.remote_exec("""
import sys ; sys.path.insert(0, %r)
""" % os.path.dirname(str(source))).waitclose()
return
if (spec, source) in self._rsynced_specs:
return
def finished():
if notify:
notify("rsyncrootready", spec, source)
rsync.add_target_host(gateway, finished=finished)
self._rsynced_specs.add((spec, source))
self.config.hook.pytest_xdist_rsyncstart(
source=source,
gateways=[gateway],
)
rsync.send()
self.config.hook.pytest_xdist_rsyncfinish(
source=source,
gateways=[gateway],
)
class HostRSync(execnet.RSync):
""" RSyncer that filters out common files
@@ -144,14 +152,12 @@ class HostRSync(execnet.RSync):
def filter(self, path):
path = py.path.local(path)
if not path.ext in ('.pyc', '.pyo'):
if not path.basename.endswith('~'):
if path.check(dotfile=0):
for x in self._ignores:
if path == x:
break
else:
return True
for x in self._ignores:
x = getattr(x, 'strpath', x)
if fnmatch.fnmatch(path.basename, x) or fnmatch.fnmatch(path.strpath, x):
return False
else:
return True
def add_target_host(self, gateway, finished=None):
remotepath = os.path.basename(self._sourcedir)
@@ -230,8 +236,8 @@ class SlaveController(object):
self.gateway.exit()
#del self.gateway
def send_runtest(self, nodeid):
self.sendcommand("runtests", ids=[nodeid])
def send_runtest_some(self, indices):
self.sendcommand("runtests", indices=indices)
def send_runtest_all(self):
self.sendcommand("runtests_all",)
@@ -278,10 +284,13 @@ class SlaveController(object):
self._down = True
self.slaveoutput = kwargs['slaveoutput']
self.notify_inproc("slavefinished", node=self)
#elif eventname == "logstart":
# self.notify_inproc(eventname, node=self, **kwargs)
elif eventname == "logstart":
self.notify_inproc(eventname, node=self, **kwargs)
elif eventname in ("testreport", "collectreport", "teardownreport"):
item_index = kwargs.pop("item_index", None)
rep = unserialize_report(eventname, kwargs['data'])
if item_index is not None:
rep.item_index = item_index
self.notify_inproc(eventname, node=self, rep=rep)
elif eventname == "collectionfinish":
self.notify_inproc(eventname, node=self, ids=kwargs['ids'])
@@ -296,8 +305,7 @@ class SlaveController(object):
self.config.pluginmanager.notify_exception(excinfo)
def unserialize_report(name, reportdict):
d = reportdict
if name == "testreport":
return runner.TestReport(**d)
return runner.TestReport(**reportdict)
elif name == "collectreport":
return runner.CollectReport(**d)
return runner.CollectReport(**reportdict)