Compare commits

..

21 Commits
1.10 ... 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
18 changed files with 725 additions and 252 deletions

View File

@@ -25,4 +25,6 @@ bin/
pytest_xdist.egg-info pytest_xdist.egg-info
issue/ issue/
3rdparty/ 3rdparty/
pytestdebug.log
.tox .tox
.cache

View File

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

View File

@@ -1,3 +1,19 @@
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
------------------------- -------------------------

View File

@@ -2,7 +2,7 @@ from setuptools import setup
setup( setup(
name="pytest-xdist", name="pytest-xdist",
version='1.10', version='1.11',
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',
@@ -13,7 +13,7 @@ setup(
packages = ['xdist'], packages = ['xdist'],
entry_points = {'pytest11': ['xdist = xdist.plugin'],}, entry_points = {'pytest11': ['xdist = xdist.plugin'],},
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'],
classifiers=[ classifiers=[
'Development Status :: 5 - Production/Stable', 'Development Status :: 5 - Production/Stable',
'Intended Audience :: Developers', 'Intended Audience :: Developers',

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):
@@ -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
@@ -250,7 +251,7 @@ 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)
@@ -300,7 +301,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 +312,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 +323,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 +352,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 +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): def test_session_hooks(testdir):
testdir.makeconftest(""" testdir.makeconftest("""
import sys import sys
@@ -459,3 +460,83 @@ def test_issue34_pluginloading_in_subprocess(testdir):
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*",
])
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): 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)
for spec in config.getvalueorskip("gspecs")]
# 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,25 @@ class TestLoadScheduling:
crashitem = sched.remove_node(node) crashitem = sched.remove_node(node)
assert crashitem == collection[0] 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: class TestDistReporter:
@@ -179,7 +198,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 +221,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

@@ -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,6 +1,7 @@
import py import py
import os import pytest
import execnet import execnet
from xdist import slavemanage
from xdist.slavemanage import HostRSync, NodeManager from xdist.slavemanage import HostRSync, NodeManager
pytest_plugins = "pytester", pytest_plugins = "pytester",
@@ -24,6 +25,14 @@ def pytest_funcarg__mysetup(request):
request.getfuncargvalue("_pytest") request.getfuncargvalue("_pytest")
return mysetup(request) 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: class TestNodeManagerPopen:
def test_popen_no_default_chdir(self, config): def test_popen_no_default_chdir(self, config):
gm = NodeManager(config, ["popen"]) gm = NodeManager(config, ["popen"])
@@ -36,9 +45,10 @@ 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, _pytest, 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 +61,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 +75,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 +98,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
@@ -108,7 +123,7 @@ class TestHRSync:
return mysetup(request) 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")
@@ -139,7 +154,7 @@ class TestNodeManager:
@py.test.mark.xfail @py.test.mark.xfail
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 +167,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 +179,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 +187,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 +199,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 +220,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
@@ -238,5 +249,3 @@ class TestNodeManager:
"--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

24
tox.ini
View File

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

View File

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

View File

@@ -1,4 +1,3 @@
import sys
import difflib import difflib
import pytest import pytest
@@ -10,35 +9,95 @@ 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
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.node2collection[node] = list(collection)
self.node2pending[node] = [] self.node2pending[node] = []
if len(self.node2pending) >= self.numnodes: if len(self.node2collection) >= self.numnodes:
self.collection_is_completed = True 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): def remove_item(self, node, item_index, duration=0):
self.node2pending[node].remove(item_index) self.node2pending[node].remove(item_index)
@@ -49,129 +108,302 @@ 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:
continue
if not pending:
pending[:] = range(len(self.node2collection[node])) pending[:] = range(len(self.node2collection[node]))
node.send_runtest_all()
else:
node.send_runtest_some(pending)
self._started.append(node)
class LoadScheduling: class LoadScheduling:
"""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): def __init__(self, numnodes, log=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
@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
if self.pending:
return False
for pending in self.node2pending.values():
if len(pending) >= 2:
return False return False
#for items in self.node2pending.values():
# if items:
# 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 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
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 +417,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")
@@ -202,6 +447,7 @@ class DSession:
self.maxfail = config.getvalue("maxfail") self.maxfail = config.getvalue("maxfail")
self.queue = queue.Queue() self.queue = queue.Queue()
self._failed_collection_errors = {} self._failed_collection_errors = {}
self._active_nodes = set()
try: try:
self.terminal = config.pluginmanager.getplugin("terminalreporter") self.terminal = config.pluginmanager.getplugin("terminalreporter")
except KeyError: except KeyError:
@@ -210,17 +456,32 @@ 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)
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()
@@ -239,7 +500,6 @@ class DSession:
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 +507,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 +528,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
if node in self.sched.nodes:
crashitem = self.sched.remove_node(node) crashitem = self.sched.remove_node(node)
#assert not crashitem, (crashitem, node) assert not crashitem, (crashitem, node)
if self.shuttingdown and not self.sched.hasnodes(): self._active_nodes.remove(node)
self.session_finished = True
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,31 +570,44 @@ 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) self.report_line("Replacing failed node %s" % node.gateway.id)
if not self.sched.hasnodes(): self._clone_node(node)
self.session_finished = True 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
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"):
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.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
@@ -328,9 +615,25 @@ class DSession:
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 +652,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

@@ -87,7 +87,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 +109,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)

View File

@@ -53,14 +53,14 @@ def pytest_addhooks(pluginmanager):
def pytest_cmdline_main(config): def pytest_cmdline_main(config):
check_options(config) check_options(config)
if config.getvalue("looponfail"): if config.getoption("looponfail"):
from xdist.looponfail import looponfail_main from xdist.looponfail import looponfail_main
looponfail_main(config) looponfail_main(config)
return 2 # looponfail only can get stop with ctrl-C anyway return 2 # looponfail only can get stop with ctrl-C anyway
def pytest_configure(config, __multicall__): def pytest_configure(config, __multicall__):
__multicall__.execute() __multicall__.execute()
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")
@@ -118,10 +118,14 @@ def forked_run_report(item):
def report_process_crash(item, result): def report_process_crash(item, result):
path, lineno = item._getfslineno() path, lineno = item._getfslineno()
info = "%s:%s: running the test CRASHED with signal %d" %( info = ("%s:%s: running the test CRASHED with signal %d" %
path, lineno, result.signal) (path, lineno, result.signal))
from _pytest import runner from _pytest import runner
call = runner.CallInfo(lambda: 0/0, "???") call = runner.CallInfo(lambda: 0/0, "???")
call.excinfo = info call.excinfo = info
rep = runner.pytest_runtest_makereport(item, call) 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 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))
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() node.setup()
self.trace("started node %r" % node) 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,38 +108,35 @@ 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()
gateways = []
for gateway in self.group:
spec = gateway.spec spec = gateway.spec
if spec.popen and not spec.chdir: if spec.popen and not spec.chdir:
# XXX this assumes that sources are python-packages # XXX This assumes that sources are python-packages
# and that adding the basedir does not hurt # and that adding the basedir does not hurt.
gateway.remote_exec(""" gateway.remote_exec("""
import sys ; sys.path.insert(0, %r) import sys ; sys.path.insert(0, %r)
""" % os.path.dirname(str(source))).waitclose() """ % os.path.dirname(str(source))).waitclose()
continue return
if spec not in seen: if (spec, source) in self._rsynced_specs:
return
def finished(): def finished():
if notify: if notify:
notify("rsyncrootready", spec, source) notify("rsyncrootready", spec, source)
rsync.add_target_host(gateway, finished=finished) rsync.add_target_host(gateway, finished=finished)
seen.add(spec) self._rsynced_specs.add((spec, source))
gateways.append(gateway)
if seen:
self.config.hook.pytest_xdist_rsyncstart( self.config.hook.pytest_xdist_rsyncstart(
source=source, source=source,
gateways=gateways, gateways=[gateway],
) )
rsync.send() rsync.send()
self.config.hook.pytest_xdist_rsyncfinish( self.config.hook.pytest_xdist_rsyncfinish(
source=source, source=source,
gateways=gateways, gateways=[gateway],
) )
class HostRSync(execnet.RSync): class HostRSync(execnet.RSync):