diff --git a/CHANGELOG b/CHANGELOG index 04c92db..b446b70 100644 --- a/CHANGELOG +++ b/CHANGELOG @@ -1,6 +1,10 @@ 1.5a1 ------------------------- + +- major internal refactoring to match the py-1.4 event refactoring + - perform test collection always at slave side instead of at the master + - make python2/python3 bridging work, remove usage of pickling - remove all trailing whitespace from source 1.4 diff --git a/testing/__init__.py b/testing/__init__.py deleted file mode 100644 index 792d600..0000000 --- a/testing/__init__.py +++ /dev/null @@ -1 +0,0 @@ -# diff --git a/testing/acceptance_test.py b/testing/acceptance_test.py index 43f12e9..9a1eac1 100644 --- a/testing/acceptance_test.py +++ b/testing/acceptance_test.py @@ -70,6 +70,18 @@ class TestDistribution: "*1 failed*", ]) + def test_basetemp_in_subprocesses(self, testdir): + p1 = testdir.makepyfile(""" + def test_send(pytestconfig): + bt = pytestconfig.getbasetemp() + assert bt.basename.startswith("popen-") + """) + result = testdir.runpytest(p1, "-n1") + assert result.ret == 0 + result.stdout.fnmatch_lines([ + "*1 passed*", + ]) + def test_dist_conftest_specified(self, testdir): p1 = testdir.makepyfile(""" import py @@ -143,28 +155,6 @@ class TestDistribution: ]) assert dest.join(subdir.basename).check(dir=1) - def test_dist_each(self, testdir): - interpreters = [] - for name in ("python2.4", "python2.5"): - interp = py.path.local.sysfind(name) - if interp is None: - py.test.skip("%s not found" % name) - interpreters.append(interp) - - testdir.makepyfile(__init__="", test_one=""" - import sys - def test_hello(): - print("%s...%s" % sys.version_info[:2]) - assert 0 - """) - args = ["--dist=each", "-v"] - args += ["--tx", "popen//python=%s" % interpreters[0]] - args += ["--tx", "popen//python=%s" % interpreters[1]] - result = testdir.runpytest(*args) - s = result.stdout.str() - assert "2.4" in s - assert "2.5" in s - def test_data_exchange(self, testdir): c1 = testdir.makeconftest(""" # This hook only called on master. @@ -236,6 +226,38 @@ class TestDistribution: child.close() #assert ret == 2 +class TestDistEach: + def test_simple(self, testdir): + testdir.makepyfile(""" + def test_hello(): + pass + """) + result = testdir.runpytest("--debug", "--dist=each", "--tx=2*popen") + assert not result.ret + result.stdout.fnmatch_lines(["*2 pass*"]) + + def test_simple_diffoutput(self, testdir): + interpreters = [] + for name in ("python2.5", "python2.6"): + interp = py.path.local.sysfind(name) + if interp is None: + py.test.skip("%s not found" % name) + interpreters.append(interp) + + testdir.makepyfile(__init__="", test_one=""" + import sys + def test_hello(): + print("%s...%s" % sys.version_info[:2]) + assert 0 + """) + args = ["--dist=each", "-v"] + args += ["--tx", "popen//python=%s" % interpreters[0]] + args += ["--tx", "popen//python=%s" % interpreters[1]] + result = testdir.runpytest(*args) + s = result.stdout.str() + assert "2...5" in s + assert "2...6" in s + class TestTerminalReporting: def test_pass_skip_fail(self, testdir): p = testdir.makepyfile(""" @@ -345,12 +367,11 @@ def test_funcarg_teardown_failure(testdir): def test_hello(myarg): pass """) - result = testdir.runpytest("--debug", p, "-n1") + result = testdir.runpytest("--debug", p) # , "-n1") result.stdout.fnmatch_lines([ "*ValueError*42*", "*1 passed*1 error*", ]) - py.test.xfail("fix exitstatus handling") assert result.ret def test_crashing_item(testdir): diff --git a/testing/test_dsession.py b/testing/test_dsession.py index df6c437..c5c68b9 100644 --- a/testing/test_dsession.py +++ b/testing/test_dsession.py @@ -1,4 +1,4 @@ -from xdist.dsession import DSession, LoadScheduling +from xdist.dsession import DSession, LoadScheduling, EachScheduling from py._test import session as outcome import py import execnet @@ -19,6 +19,9 @@ class MockNode: def send_runtest(self, nodeid): self.sent.append(nodeid) + def send_runtest_all(self): + self.sent.append("ALL") + def sendlist(self, items): self.sent.extend(items) @@ -29,6 +32,46 @@ def dumpqueue(queue): while queue.qsize(): print(queue.get()) +class TestEachScheduling: + def test_schedule_load_simple(self): + node1 = MockNode() + node2 = MockNode() + sched = EachScheduling(2) + sched.addnode(node1) + sched.addnode(node2) + collection = ["a.py::test_1", ] + assert not sched.collection_is_completed + sched.addnode_collection(node1, collection) + assert not sched.collection_is_completed + sched.addnode_collection(node2, collection) + assert sched.collection_is_completed + assert sched.node2collection[node1] == collection + assert sched.node2collection[node2] == collection + sched.init_distribute() + assert not sched.tests_finished() + assert node1.sent == ['ALL'] + assert node2.sent == ['ALL'] + sched.remove_item(node1, collection[0]) + assert not sched.tests_finished() + sched.remove_item(node2, collection[0]) + assert sched.tests_finished() + + def test_schedule_remove_node(self): + node1 = MockNode() + sched = EachScheduling(1) + sched.addnode(node1) + collection = ["a.py::test_1", ] + assert not sched.collection_is_completed + sched.addnode_collection(node1, collection) + assert sched.collection_is_completed + assert sched.node2collection[node1] == collection + sched.init_distribute() + assert not sched.tests_finished() + crashitem = sched.remove_node(node1) + assert crashitem + assert sched.tests_finished() + assert not sched.hasnodes() + class TestLoadScheduling: def test_schedule_load_simple(self): node1 = MockNode() @@ -45,17 +88,17 @@ class TestLoadScheduling: assert sched.node2collection[node1] == collection assert sched.node2collection[node2] == collection sched.init_distribute() - assert sched.pending - sched.triggertesting() assert not sched.tests_finished() - assert node1.sent == collection[:1] - assert node2.sent == collection[1:] - sched.remove_item(node1, collection[0]) - sched.remove_item(node2, collection[1]) + 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 - def test_triggertesting_chunksize(self): + def test_init_distribute_chunksize(self): sched = LoadScheduling(2) node1 = MockNode() node2 = MockNode() @@ -66,7 +109,6 @@ class TestLoadScheduling: sched.addnode_collection(node1, col) sched.addnode_collection(node2, col) sched.init_distribute() - sched.triggertesting() sent1 = node1.sent sent2 = node2.sent chunkitems = col[:sched.ITEM_CHUNKSIZE] @@ -78,7 +120,6 @@ class TestLoadScheduling: for node in (node1, node2): for i in range(sched.ITEM_CHUNKSIZE): sched.remove_item(node, "xyz") - sched.triggertesting() assert not sched.pending def test_add_remove_node(self): @@ -89,320 +130,14 @@ class TestLoadScheduling: sched.addnode_collection(node, collection) assert sched.collection_is_completed sched.init_distribute() - sched.triggertesting() assert not sched.pending crashitem = sched.remove_node(node) assert crashitem == collection[0] -class TestDSession: - - #def test_collection_fails(self, testdir): - # pass - - def xxx_test_senditems_each_and_receive_with_two_nodes(self, testdir): - item = testdir.getitem("def test_func(): pass") - node1 = MockNode() - node2 = MockNode() - session = DSession(item.config) - session.addnode(node1) - session.addnode(node2) - session.senditems_each([item]) - assert session.node2pending[node1] == [item] - assert session.node2pending[node2] == [item] - assert node1 in session.item2nodes[item] - assert node2 in session.item2nodes[item] - session.removeitem(item, node1) - assert session.item2nodes[item] == [node2] - session.removeitem(item, node2) - assert not session.node2pending[node1] - assert not session.item2nodes - - def test_keyboardinterrupt(self, testdir): - item = testdir.getitem("def test_func(): pass") - session = DSession(item.config) - def raise_(timeout=None): raise KeyboardInterrupt() - session.queue.get = raise_ - exitstatus = session.loop([]) - assert exitstatus == outcome.EXIT_INTERRUPTED - - def test_internalerror(self, testdir): - item = testdir.getitem("def test_func(): pass") - session = DSession(item.config) - def raise_(): raise ValueError() - session.queue.get = raise_ - exitstatus = session.loop([]) - assert exitstatus == outcome.EXIT_INTERNALERROR - - def test_no_node_remaining_for_tests(self, testdir): - item = testdir.getitem("def test_func(): pass") - # setup a session with one node - session = DSession(item.config) - node = MockNode() - session.addnode(node) - - # setup a HostDown event - session.queueevent("pytest_testnodedown", node=node, error=None) - - loopstate = session._initloopstate([item]) - loopstate.dowork = False - session.loop_once(loopstate) - dumpqueue(session.queue) - assert loopstate.exitstatus == outcome.EXIT_NOHOSTS - - def test_removeitem_from_failing_teardown(self, testdir): - # teardown reports only come in when they signal a failure - # internal session-management should basically ignore them - # XXX probably it'S best to invent a new error hook for - # teardown/setup related failures - modcol = testdir.getmodulecol(""" - def test_one(): - pass - def teardown_function(function): - assert 0 - """) - item1, = modcol.collect() - - # setup a session with two nodes - session = DSession(item1.config) - node1, node2 = MockNode(), MockNode() - session.addnode(node1) - session.addnode(node2) - - # have one test pending for a node that goes down - session.senditems_each([item1]) - nodes = session.item2nodes[item1] - class rep: - failed = True - item = item1 - node = nodes[0] - when = "call" - session.queueevent("pytest_runtest_logreport", report=rep) - reprec = testdir.getreportrecorder(session) - print(session.item2nodes) - loopstate = session._initloopstate([]) - assert len(session.item2nodes[item1]) == 2 - session.loop_once(loopstate) - assert len(session.item2nodes[item1]) == 1 - rep.when = "teardown" - session.queueevent("pytest_runtest_logreport", report=rep) - session.loop_once(loopstate) - assert len(session.item2nodes[item1]) == 1 - - def test_testnodeready_adds_to_available(self, testdir): - item = testdir.getitem("def test_func(): pass") - # setup a session with two nodes - session = DSession(item.config) - node1 = MockNode() - session.queueevent("pytest_testnodeready", node=node1) - loopstate = session._initloopstate([item]) - loopstate.dowork = False - assert len(session.node2pending) == 0 - session.loop_once(loopstate) - assert len(session.node2pending) == 1 - - def runthrough(self, item, excinfo=None): - session = DSession(item.config) - node = MockNode() - session.addnode(node) - loopstate = session._initloopstate([item]) - - session.queueevent(None) - session.loop_once(loopstate) - - assert node.sent == [item] - ev = run(item, node, excinfo=excinfo) - session.queueevent("pytest_runtest_logreport", report=ev) - session.loop_once(loopstate) - assert loopstate.shuttingdown - session.queueevent("pytest_testnodedown", node=node, error=None) - session.loop_once(loopstate) - dumpqueue(session.queue) - return session, loopstate.exitstatus - - def test_exit_completed_tests_ok(self, testdir): - item = testdir.getitem("def test_func(): pass") - session, exitstatus = self.runthrough(item) - assert exitstatus == outcome.EXIT_OK - - def test_exit_completed_tests_fail(self, testdir): - item = testdir.getitem("def test_func(): 0/0") - session, exitstatus = self.runthrough(item, excinfo="fail") - assert exitstatus == outcome.EXIT_TESTSFAILED - - def test_exit_on_first_failing(self, testdir): - modcol = testdir.getmodulecol(""" - def test_fail(): - assert 0 - def test_pass(): - pass - """) - modcol.config.option.maxfail = 1 - session = DSession(modcol.config) - node = MockNode() - session.addnode(node) - items = modcol.config.hook.pytest_make_collect_report(collector=modcol).result - - # trigger testing - this sends tests to the node - session.triggertesting(items) - - # run tests ourselves and produce reports - ev1 = run(items[0], node, "fail") - ev2 = run(items[1], node, None) - session.queueevent("pytest_runtest_logreport", report=ev1) - session.queueevent("pytest_runtest_logreport", report=ev2) - # now call the loop - loopstate = session._initloopstate(items) - py.test.raises(session.Interrupted, "session.loop_once(loopstate)") - assert loopstate.testsfailed - #assert loopstate.shuttingdown - - def test_maxfail(self, testdir): - modcol = testdir.getmodulecol(""" - def test_fail1(): - assert 0 - def test_fail2(): - assert 0 - def test_pass(): - pass - """) - modcol.config.option.maxfail = 2 - session = DSession(modcol.config) - node = MockNode() - session.addnode(node) - items = modcol.config.hook.pytest_make_collect_report(collector=modcol).result - - # trigger testing - this sends tests to the node - session.triggertesting(items) - - # run tests ourselves and produce reports - ev1 = run(items[0], node, "fail") - ev2 = run(items[1], node, "fail") - session.queueevent("pytest_runtest_logreport", report=ev1) # a failing one - session.queueevent("pytest_runtest_logreport", report=ev2) - # now call the loop - loopstate = session._initloopstate(items) - try: - session.loop_once(loopstate) - except session.Interrupted: - py.test.fail("raised Interrupted but shouildn't") - py.test.raises(session.Interrupted, "session.loop_once(loopstate)") - assert loopstate.testsfailed - #assert loopstate.shuttingdown - - def test_shuttingdown_filters(self, testdir): - item = testdir.getitem("def test_func(): pass") - session = DSession(item.config) - node = MockNode() - session.addnode(node) - loopstate = session._initloopstate([]) - loopstate.shuttingdown = True - reprec = testdir.getreportrecorder(session) - session.queueevent("pytest_runtest_logreport", report=run(item, node)) - session.loop_once(loopstate) - assert not reprec.getcalls("pytest_testnodedown") - session.queueevent("pytest_testnodedown", node=node, error=None) - session.loop_once(loopstate) - assert reprec.getcall('pytest_testnodedown').node == node - - def test_filteritems(self, testdir): - modcol = testdir.getmodulecol(""" - def test_fail(): - assert 0 - def test_pass(): - pass - """) - session = DSession(modcol.config) - - modcol.config.option.keyword = "nothing" - dsel = session.filteritems([modcol]) - assert dsel == [modcol] - items = modcol.collect() - hookrecorder = testdir.getreportrecorder(session).hookrecorder - remaining = session.filteritems(items) - assert remaining == [] - - event = hookrecorder.getcalls("pytest_deselected")[-1] - assert event.items == items - - modcol.config.option.keyword = "test_fail" - remaining = session.filteritems(items) - assert remaining == [items[0]] - - event = hookrecorder.getcalls("pytest_deselected")[-1] - assert event.items == [items[1]] - - def test_testnodedown_shutdown_after_completion(self, testdir): - item = testdir.getitem("def test_func(): pass") - session = DSession(item.config) - - node = MockNode() - session.addnode(node) - session.senditems_load([item]) - session.queueevent("pytest_runtest_logreport", report=run(item, node)) - loopstate = session._initloopstate([]) - session.loop_once(loopstate) - assert node._shutdown is True - assert loopstate.exitstatus is None, "loop did not wait for testnodedown" - assert loopstate.shuttingdown - session.queueevent("pytest_testnodedown", node=node, error=None) - session.loop_once(loopstate) - assert loopstate.exitstatus == 0 - - def test_nopending_but_collection_remains(self, testdir): - modcol = testdir.getmodulecol(""" - def test_fail(): - assert 0 - def test_pass(): - pass - """) - session = DSession(modcol.config) - node = MockNode() - session.addnode(node) - - colreport = modcol.config.hook.pytest_make_collect_report(collector=modcol) - item1, item2 = colreport.result - session.senditems_load([item1]) - # node2pending will become empty when the loop sees the report - rep = run(item1, node) - session.queueevent("pytest_runtest_logreport", report=run(item1, node)) - - # but we have a collection pending - session.queueevent("pytest_collectreport", report=colreport) - - loopstate = session._initloopstate([]) - session.loop_once(loopstate) - assert loopstate.exitstatus is None, "loop did not care for collection report" - assert not loopstate.colitems - session.loop_once(loopstate) - assert loopstate.colitems == colreport.result - assert loopstate.exitstatus is None, "loop did not care for colitems" - - def test_dist_some_tests(self, testdir): - p1 = testdir.makepyfile(test_one=""" - def test_1(): - pass - def test_x(): - import py - py.test.skip("aaa") - def test_fail(): - assert 0 - """) - config = testdir.parseconfig('-d', p1, '--tx=popen') - dsession = DSession(config) - hookrecorder = testdir.getreportrecorder(config).hookrecorder - dsession.main([config.getnode(p1)]) - rep = hookrecorder.popcall("pytest_runtest_logreport").report - assert rep.passed - rep = hookrecorder.popcall("pytest_runtest_logreport").report - assert rep.skipped - rep = hookrecorder.popcall("pytest_runtest_logreport").report - assert rep.failed - # see that the node is really down - node = hookrecorder.popcall("pytest_testnodedown").node - assert node.gateway.spec.popen - #XXX eq.geteventargs("pytest_sessionfinish") +class TestDistReporter: + @py.test.mark.xfail def test_rsync_printing(self, testdir, linecomp): config = testdir.parseconfig() from py._plugin.pytest_terminal import TerminalReporter diff --git a/testing/test_remote.py b/testing/test_remote.py index 9082dc1..3b3a284 100644 --- a/testing/test_remote.py +++ b/testing/test_remote.py @@ -174,6 +174,34 @@ class TestSlaveInteractor: print ev.kwargs assert not ev.kwargs['ids'] + def test_runtests_all(self, slave): + p = slave.testdir.makepyfile(""" + def test_func(): pass + def test_func2(): pass + """) + slave.setup() + ev = slave.popevent() + assert ev.name == "slaveready" + ev = slave.popevent() + assert ev.name == "collectionstart" + assert not ev.kwargs + ev = slave.popevent("collectionfinish") + ids = ev.kwargs['ids'] + assert len(ids) == 2 + slave.sendcommand("runtests_all", ) + ev = slave.popevent("testreport") + assert ev.name == "testreport" + rep = unserialize_report(ev.name, ev.kwargs['data']) + assert rep.nodeid.endswith("::test_func") + ev = slave.popevent("testreport") + assert ev.name == "testreport" + rep = unserialize_report(ev.name, ev.kwargs['data']) + assert rep.nodeid.endswith("::test_func2") + assert rep.passed + slave.sendcommand("shutdown") + ev = slave.popevent("slavefinished") + assert 'slaveoutput' in ev.kwargs + def test_happy_run_events_converted(self, testdir, slave): py.test.xfail("implement a simple test for event production") assert not slave.use_callback diff --git a/testing/test_txnode.py b/testing/test_txnode.py deleted file mode 100644 index 17c92f1..0000000 --- a/testing/test_txnode.py +++ /dev/null @@ -1,172 +0,0 @@ - -import py -import execnet -from xdist.txnode import TXNode -queue = py.builtin._tryimport("queue", "Queue") -Queue = queue.Queue - -class EventQueue: - def __init__(self, registry, queue=None): - if queue is None: - queue = Queue() - self.queue = queue - registry.register(self) - - def geteventargs(self, eventname, timeout=10.0): - events = [] - while 1: - try: - eventcall = self.queue.get(timeout=timeout) - except queue.Empty: - #print "node channel", self.node.channel - #print "remoteerror", self.node.channel._getremoteerror() - py.builtin.print_("seen events", events) - raise IOError("did not see %r events" % (eventname)) - else: - name, args, kwargs = eventcall - assert isinstance(name, str) - if name == eventname: - if args: - return args - return kwargs - events.append(name) - if name == "pytest_internalerror": - py.builtin.print_(str(kwargs["excrepr"])) - -class MySetup: - def __init__(self, request): - self.id = 0 - self.request = request - - def geteventargs(self, eventname, timeout=10.0): - eq = EventQueue(self.config.pluginmanager, self.queue) - return eq.geteventargs(eventname, timeout=timeout) - - def makenode(self, config=None, xspec="popen"): - if config is None: - testdir = self.request.getfuncargvalue("testdir") - config = testdir.reparseconfig([]) - self.config = config - self.queue = Queue() - self.xspec = execnet.XSpec(xspec) - self.gateway = execnet.makegateway(self.xspec) - self.id += 1 - self.gateway.id = str(self.id) - self.nodemanager = None - self.node = TXNode(self.nodemanager, self.gateway, self.config, putevent=self.queue.put) - assert not self.node.channel.isclosed() - return self.node - - def xfinalize(self): - if hasattr(self, 'node'): - gw = self.node.gateway - py.builtin.print_("exiting:", gw) - gw.exit() - -def pytest_funcarg__mysetup(request): - mysetup = MySetup(request) - #pyfuncitem.addfinalizer(mysetup.finalize) - return mysetup - -def test_node_hash_equality(mysetup): - node = mysetup.makenode() - node2 = mysetup.makenode() - assert node != node2 - assert node == node - assert not (node != node) - -class TestMasterSlaveConnection: - def test_crash_invalid_item(self, mysetup): - node = mysetup.makenode() - node.send(123) # invalid item - kwargs = mysetup.geteventargs("pytest_testnodedown") - assert kwargs['node'] is node - #assert isinstance(kwargs['error'], execnet.RemoteError) - - def test_crash_killed(self, testdir, mysetup): - if not hasattr(py.std.os, 'kill'): - py.test.skip("no os.kill") - item = testdir.getitem(""" - def test_func(): - import os - os.kill(os.getpid(), 9) - """) - node = mysetup.makenode(item.config) - node.send(item) - kwargs = mysetup.geteventargs("pytest_testnodedown") - assert kwargs['node'] is node - assert "Not properly terminated" in str(kwargs['error']) - - def test_node_down(self, mysetup): - node = mysetup.makenode() - node.shutdown() - kwargs = mysetup.geteventargs("pytest_testnodedown") - assert kwargs['node'] is node - assert not kwargs['error'] - node.callback(node.ENDMARK) - excinfo = py.test.raises(IOError, - "mysetup.geteventargs('testnodedown', timeout=0.01)") - - def test_send_on_closed_channel(self, testdir, mysetup): - item = testdir.getitem("def test_func(): pass") - node = mysetup.makenode(item.config) - node.channel.close() - py.test.raises(IOError, "node.send(item)") - #ev = self.getcalls(pytest_internalerror) - #assert ev.excinfo.errisinstance(IOError) - - def test_send_one(self, testdir, mysetup): - item = testdir.getitem("def test_func(): pass") - node = mysetup.makenode(item.config) - node.send(item) - kwargs = mysetup.geteventargs("pytest_runtest_logreport") - rep = kwargs['report'] - assert rep.passed - py.builtin.print_(rep) - assert rep.item == item - - def test_send_some(self, testdir, mysetup): - items = testdir.getitems(""" - def test_pass(): - pass - def test_fail(): - assert 0 - def test_skip(): - import py - py.test.skip("x") - """) - node = mysetup.makenode(items[0].config) - for item in items: - node.send(item) - for outcome in "passed failed skipped".split(): - kwargs = mysetup.geteventargs("pytest_runtest_logreport") - report = kwargs['report'] - assert getattr(report, outcome) - - node.sendlist(items) - for outcome in "passed failed skipped".split(): - rep = mysetup.geteventargs("pytest_runtest_logreport")['report'] - assert getattr(rep, outcome) - - def test_send_one_with_env(self, testdir, mysetup, monkeypatch): - if execnet.XSpec("popen").env is None: - py.test.skip("requires execnet 1.0.7 or above") - monkeypatch.delenv('ENV1', raising=False) - monkeypatch.delenv('ENV2', raising=False) - monkeypatch.setenv('ENV3', 'var3') - - item = testdir.getitem(""" - def test_func(): - import os - # ENV1, ENV2 set by xspec; ENV3 inherited from parent process - assert os.getenv('ENV2') == 'var2' - assert os.getenv('ENV1') == 'var1' - assert os.getenv('ENV3') == 'var3' - """) - node = mysetup.makenode(item.config, - xspec="popen//env:ENV1=var1//env:ENV2=var2") - node.send(item) - kwargs = mysetup.geteventargs("pytest_runtest_logreport") - rep = kwargs['report'] - assert rep.passed - diff --git a/xdist/dsession.py b/xdist/dsession.py index 128a606..416d150 100644 --- a/xdist/dsession.py +++ b/xdist/dsession.py @@ -1,17 +1,59 @@ import py import sys from xdist.slavemanage import NodeManager -from py._test import session queue = py.builtin._tryimport('queue', 'Queue') -def dsession_main(config): - config.pluginmanager.do_configure(config) - session = DSession(config) - trdist = TerminalDistReporter(config) - config.pluginmanager.register(trdist, "terminaldistreporter") - exitcode = session.main() - config.pluginmanager.do_unconfigure(config) - return exitcode +class EachScheduling: + + def __init__(self, numnodes, log=None): + self.numnodes = numnodes + self.node2collection = {} + self.node2pending = {} + if log is None: + self.log = py.log.Producer("eachsched") + else: + self.log = log.loadsched + self.collection_is_completed = False + + def hasnodes(self): + return bool(self.node2pending) + + def addnode(self, node): + self.node2collection[node] = None + + def tests_finished(self): + if not self.collection_is_completed: + return False + for items in self.node2pending.values(): + if items: + 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 + + def remove_item(self, node, item): + self.node2pending[node].remove(item) + + 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? + return crashitem + + def init_distribute(self): + assert self.collection_is_completed + for node, pending in self.node2pending.items(): + node.send_runtest_all() + pending[:] = self.node2collection[node] class LoadScheduling: LOAD_THRESHOLD_NEWITEMS = 5 @@ -92,29 +134,20 @@ class LoadScheduling: for node, collection in self.node2collection.items(): assert collection == col self.pending = col - - def triggertesting(self): - if not self.pending: + if not col: return - available = [] - for node, pending in self.node2pending.items(): - if len(pending) < self.LOAD_THRESHOLD_NEWITEMS: - available.append((node, pending)) - num_available = len(available) - if num_available: - 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) - #item.ihook.pytest_itemstart(item=item, node=node) - pending.append(item) - if i >= max_one_round: - break - del self.pending[:i+1] - if self.pending: - self.log("triggertesting remaining:", len(self.pending)) + 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] class Interrupted(KeyboardInterrupt): """ signals an immediate interruption. """ @@ -138,23 +171,59 @@ class DSession: if self.terminal: self.terminal.write_line(line) - def pytest_gwmanage_rsyncstart(self, source, gateways): - targets = ",".join([gw.id for gw in gateways]) - msg = "[%s] rsyncing: %s" %(targets, source) - self.report_line(msg) + def pytest_sessionstart(self, session, __multicall__): + #print "remaining multicall methods", __multicall__.methods + if not self.config.getvalue("verbose"): + self.report_line("instantiating gateways (use -v for details): %s" % + ",".join(self.config.option.tx)) + self.nodemanager = NodeManager(self.config) + self.nodemanager.setup_nodes(putevent=self.queue.put) - #def pytest_gwmanage_rsyncfinish(self, source, gateways): - # targets = ", ".join(["[%s]" % gw.id for gw in gateways]) - # self.write_line("rsyncfinish: %s -> %s" %(source, targets)) + def pytest_sessionfinish(self, session): + """ teardown any resources after a test run. """ + self.nodemanager.teardown_nodes() - def main(self): - self.config.hook.pytest_sessionstart(session=self) - self.setup() - exitstatus = self.loop() - self.teardown() - self.config.hook.pytest_sessionfinish(session=self, - exitstatus=exitstatus,) - return exitstatus + def pytest_perform_collection(self, __multicall__): + # prohibit collection of test items in master process + __multicall__.methods[:] = [] + + def pytest_runtest_mainloop(self): + numnodes = len(self.nodemanager.gwmanager.specs) + dist = self.config.getvalue("dist") + if dist == "load": + self.sched = LoadScheduling(numnodes, log=self.log) + elif dist == "each": + self.sched = EachScheduling(numnodes, log=self.log) + else: + assert 0, dist + self.shouldstop = False + self.session_finished = False + while not self.session_finished: + self.loop_once() + if self.shouldstop: + raise Interrupted(str(self.shouldstop)) + return True + + def loop_once(self): + """ process one callback from one of the slaves. """ + while 1: + try: + eventcall = self.queue.get(timeout=2.0) + break + except queue.Empty: + continue + callname, kwargs = eventcall + assert callname, kwargs + method = "slave_" + callname + call = getattr(self, method) + self.log("calling method: %s(**%s)" % (method, kwargs)) + call(**kwargs) + if self.sched.tests_finished(): + self.triggershutdown() + + # + # callbacks for processing events from slaves + # def slave_slaveready(self, node, slaveinfo): node.slaveinfo = slaveinfo @@ -192,7 +261,6 @@ class DSession: if self.sched.collection_is_completed: self.sched.init_distribute() - self.sched.triggertesting() def slave_logstart(self, node, nodeid, location): self.config.hook.pytest_runtest_logstart( @@ -222,45 +290,6 @@ class DSession: self.shouldstop = "stopping after %d failures" % ( self.countfailures) - def loop(self): - numnodes = len(self.nodemanager.gwmanager.specs) - self.sched = LoadScheduling(numnodes, log=self.log) - self.shouldstop = False - self.session_finished = False - exitstatus = 0 - try: - while not self.session_finished: - self.loop_once() - if self.shouldstop: - raise Interrupted(str(self.shouldstop)) - except KeyboardInterrupt: - excinfo = py.code.ExceptionInfo() - self.config.hook.pytest_keyboard_interrupt(excinfo=excinfo) - exitstatus = session.EXIT_INTERRUPTED - except: - self.config.pluginmanager.notify_exception() - exitstatus = session.EXIT_INTERNALERROR - #self.config.pluginmanager.unregister(loopstate) - if exitstatus == 0 and self.countfailures: - exitstatus = session.EXIT_TESTSFAILED - return exitstatus - - def loop_once(self): - while 1: - try: - eventcall = self.queue.get(timeout=2.0) - break - except queue.Empty: - continue - callname, kwargs = eventcall - assert callname, kwargs - method = "slave_" + callname - call = getattr(self, method) - self.log("calling method: %s(**%s)" % (method, kwargs)) - call(**kwargs) - if self.sched.tests_finished(): - self.triggershutdown() - def triggershutdown(self): self.log("triggering shutdown") self.shuttingdown = True @@ -277,18 +306,6 @@ class DSession: enrich_report_with_platform_data(rep, slave) self.config.hook.pytest_runtest_logreport(report=rep) - def setup(self): - """ setup any neccessary resources ahead of the test run. """ - if not self.config.getvalue("verbose"): - self.report_line("instantiating gateways (use -v for details): %s" % - ",".join(self.config.option.tx)) - self.nodemanager = NodeManager(self.config) - self.nodemanager.setup_nodes(putevent=self.queue.put) - - def teardown(self): - """ teardown any resources after a test run. """ - self.nodemanager.teardown_nodes() - class TerminalDistReporter: def __init__(self, config): self.config = config @@ -305,9 +322,9 @@ class TerminalDistReporter: gateway.id, rinfo.platform, version, rinfo.cwd)) def pytest_testnodeready(self, node): - if self.config.getvalue("debug"): + if self.config.getvalue("verbose"): d = node.slaveinfo - infoline = "[%s] -- Python %s" %( + infoline = "[%s] Python %s" %( d['id'], d['version'].replace('\n', ' -- '),) self.write_line(infoline) @@ -317,6 +334,15 @@ class TerminalDistReporter: return self.write_line("[%s] node down: %s" %(node.gateway.id, error)) + #def pytest_gwmanage_rsyncstart(self, source, gateways): + # targets = ",".join([gw.id for gw in gateways]) + # msg = "[%s] rsyncing: %s" %(targets, source) + # self.write_line(msg) + #def pytest_gwmanage_rsyncfinish(self, source, gateways): + # targets = ", ".join(["[%s]" % gw.id for gw in gateways]) + # self.write_line("rsyncfinish: %s -> %s" %(source, targets)) + + def enrich_report_with_platform_data(rep, node): rep.node = node if hasattr(rep, 'node') and rep.longrepr: diff --git a/xdist/plugin.py b/xdist/plugin.py index 6dbd72f..5179632 100644 --- a/xdist/plugin.py +++ b/xdist/plugin.py @@ -190,9 +190,16 @@ def pytest_cmdline_main(config): from xdist.looponfail import looponfail_main looponfail_main(config) return 2 # looponfail only can get stop with ctrl-C anyway - elif config.getvalue("dist") != "no": - from xdist.dsession import dsession_main - return dsession_main(config) + +def pytest_configure(config, __multicall__): + __multicall__.execute() + if config.getvalue("dist") != "no": + from xdist.dsession import DSession, TerminalDistReporter + session = DSession(config) + config.pluginmanager.register(session, "dsession") + + trdist = TerminalDistReporter(config) + config.pluginmanager.register(trdist, "terminaldistreporter") def check_options(config): if config.option.numprocesses: @@ -224,7 +231,8 @@ def forked_run_report(item): from py._plugin.pytest_runner import runtestprotocol EXITSTATUS_TESTEXIT = 4 import marshal - from xdist.remote import serialize_report, unserialize_report + from xdist.remote import serialize_report + from xdist.slavemanage import unserialize_report def runforked(): try: reports = runtestprotocol(item, log=False) @@ -236,7 +244,7 @@ def forked_run_report(item): result = ff.waitfinish() if result.retval is not None: report_dumps = marshal.loads(result.retval) - return [unserialize_report(x) for x in report_dumps] + return [unserialize_report("testreport", x) for x in report_dumps] else: if result.exitstatus == EXITSTATUS_TESTEXIT: py.test.exit("forked test item %s raised Exit" %(item,)) diff --git a/xdist/remote.py b/xdist/remote.py index a763b7c..770b4b8 100644 --- a/xdist/remote.py +++ b/xdist/remote.py @@ -55,6 +55,9 @@ class SlaveInteractor: for nodeid in ids: for item in self.collection.getbyid(nodeid): self.config.hook.pytest_runtest_protocol(item=item) + elif name == "runtests_all": + for item in self.collection.items: + self.config.hook.pytest_runtest_protocol(item=item) elif name == "shutdown": break return True @@ -65,7 +68,7 @@ class SlaveInteractor: topdir=str(collection.topdir), ids=ids) - #def pytest_runtest_logstart(self, nodeid, location): + #def pytest_runtest_logstart(self, nodeid, location, fspath): # self.sendevent("logstart", nodeid=nodeid, location=location) def pytest_runtest_logreport(self, report): diff --git a/xdist/slavemanage.py b/xdist/slavemanage.py index aa32a53..85c39da 100644 --- a/xdist/slavemanage.py +++ b/xdist/slavemanage.py @@ -233,9 +233,13 @@ class SlaveController(object): args = self.config.args if not spec.popen or spec.chdir: args = make_reltoroot(self.nodemanager.roots, args) + option_dict = vars(self.config.option) + if spec.popen: + name = "popen-%s" % self.gateway.id + option_dict['basetemp'] = str(self.config.getbasetemp().join(name)) self.config.hook.pytest_configure_node(node=self) self.channel = self.gateway.remote_exec(xdist.remote) - self.channel.send((self.slaveinput, args, vars(self.config.option))) + self.channel.send((self.slaveinput, args, option_dict)) if self.putevent: self.channel.setcallback(self.process_from_remote, endmarker=self.ENDMARK) @@ -254,6 +258,9 @@ class SlaveController(object): def send_runtest(self, nodeid): self.sendcommand("runtests", ids=[nodeid]) + def send_runtest_all(self): + self.sendcommand("runtests_all",) + def shutdown(self): if not self._down and not self.channel.isclosed(): self.sendcommand("shutdown") diff --git a/xdist/txnode.py b/xdist/txnode.py deleted file mode 100644 index acceb0c..0000000 --- a/xdist/txnode.py +++ /dev/null @@ -1,175 +0,0 @@ -""" - Manage setup, running and local representation of remote nodes/processes. -""" -import py -from py._test.session import Session - -class TXNode(object): - """ Represents a Test Execution environment in the controlling process. - - sets up a slave node through an execnet gateway - - manages sending of test-items and receival of results and events - - creates events when the remote side crashes - """ - ENDMARK = -1 - - def __init__(self, nodemanager, gateway, config, putevent): - self.nodemanager = nodemanager - self.config = config - self.putevent = putevent - self.gateway = gateway - self.slaveinput = {} - self.channel = self.setup() - self.channel.setcallback(self.callback, endmarker=self.ENDMARK) - self._down = False - - def __repr__(self): - id = self.gateway.id - status = self._down and 'true' or 'false' - return "" %(id, status) - - def notify(self, eventname, *args, **kwargs): - assert not args - self.putevent((eventname, args, kwargs)) - - def callback(self, eventcall): - """ this gets called for each object we receive from - the other side and if the channel closes. - - Note that channel callbacks run in the receiver - thread of execnet gateways - we need to - avoid raising exceptions or doing heavy work. - """ - try: - if eventcall == self.ENDMARK: - err = self.channel._getremoteerror() - if not self._down: - if not err or isinstance(err, EOFError): - err = "Not properly terminated" # lost connection? - self.notify("pytest_testnodedown", node=self, error=err) - self._down = True - return - eventname, args, kwargs = eventcall - if eventname == "slaveready": - self.notify("pytest_testnodeready", node=self) - elif eventname == "slavefinished": - self._down = True - self.slaveoutput = kwargs['slaveoutput'] - error = kwargs['error'] - self.notify("pytest_testnodedown", error=error, node=self) - elif eventname in ("pytest_runtest_logreport", - "pytest__teardown_final_logerror"): - kwargs['report'].node = self - self.notify(eventname, **kwargs) - else: - self.notify(eventname, **kwargs) - except KeyboardInterrupt: - # should not land in receiver-thread - raise - except: - excinfo = py.code.ExceptionInfo() - py.builtin.print_("!" * 20, excinfo) - self.config.pluginmanager.notify_exception(excinfo) - - def send(self, item): - assert item is not None - self.channel.send(item) - - def sendlist(self, itemlist): - self.channel.send(itemlist) - - def shutdown(self, kill=False): - if kill: - self.gateway.exit() - else: - self.channel.send(None) - - # configuring and setting up slave node - def setup(self): - basetemp = None - config = self.config - config.hook.pytest_configure_node(node=self) - if self.gateway.spec.popen: - popenbase = config.ensuretemp("popen") - basetemp = py.path.local.make_numbered_dir(prefix="slave-", - keep=0, rootdir=popenbase) - basetemp = str(basetemp) - return self.gateway.remote_exec(init_slave_session, - args=self.config.args, - option_dict=vars(self.config.option), - slaveinput={}, # XXX, - basetemp=basetemp, - nodeid=self.gateway.id, - ) - -def init_slave_session(channel, args, option_dict, - slaveinput, basetemp, nodeid): - import os, sys - #sys.path.insert(0, os.getcwd()) - from xdist.txnode import SlaveSession - import py - config = py.test.config - config.option.__dict__.update(option_dict) - config._preparse(args) - config.args = args - config.slaveinput = slaveinput - config.slaveoutput = {} - if basetemp: - config.basetemp = py.path.local(basetemp) - config.nodeid = nodeid - return SlaveSession(config, channel).dist_main() - -class SlaveSession: - def __init__(self, config, channel): - self.channel = channel - self.config = config - self.runner = self.config.pluginmanager.getplugin("pytest_runner") - config.pluginmanager.register(self, "slavesession") - - def __repr__(self): - return "<%s channel=%s>" %(self.__class__.__name__, self.channel) - - def sendevent(self, eventname, *args, **kwargs): - self.channel.send((eventname, args, kwargs)) - - def pytest_runtest_logreport(self, report): - self.sendevent("pytest_runtest_logreport", report=report) - - def pytest__teardown_final_logerror(self, report): - self.sendevent("pytest__teardown_final_logerror", report=report) - - def pytest_keyboard_interrupt(self, excinfo): - self._slaveerror = "SIGINT" - - def pytest_internalerror(self, excrepr): - self._slaveerror = "internal-error" - self.sendevent("pytest_internalerror", excrepr=excrepr) - - def dist_main(self): - self.sendevent("slaveready") - self.main(None) - error = getattr(self, '_slaveerror', None) - self.sendevent("slavefinished", error=error, - slaveoutput=self.config.slaveoutput) - - def _mainloop(self, colitems): - while 1: - task = self.channel.receive() - if task is None: - break - if isinstance(task, list): - for item in task: - self.run_single(item=item) - else: - self.run_single(item=task) - - def run_single(self, item): - call = self.runner.CallInfo(item._reraiseunpicklingproblem, when='setup') - if call.excinfo: - # likely it is not collectable here because of - # platform/import-dependency induced skips - # we fake a setup-error report with the obtained exception - # and do not care about capturing or non-runner hooks - rep = self.runner.pytest_runtest_makereport(item=item, call=call) - self.pytest_runtest_logreport(rep) - return - item.config.hook.pytest_runtest_protocol(item=item)