From f7923c692ea9448425ab9477f7b7b147948ff91d Mon Sep 17 00:00:00 2001 From: holger krekel Date: Mon, 27 Sep 2010 16:13:57 +0200 Subject: [PATCH] move code to fewer more sensible modules. --- testing/test_gwmanage.py | 121 -------------- testing/test_nodemanage.py | 120 -------------- testing/test_plugin.py | 10 +- testing/test_remote.py | 4 +- testing/test_slavemanage.py | 241 +++++++++++++++++++++++++++ xdist/dsession.py | 2 +- xdist/newhooks.py | 2 +- xdist/nodemanage.py | 101 ------------ xdist/remote.py | 180 ++++---------------- xdist/slavemanage.py | 318 ++++++++++++++++++++++++++++++++++++ 10 files changed, 598 insertions(+), 501 deletions(-) delete mode 100644 testing/test_gwmanage.py delete mode 100644 testing/test_nodemanage.py create mode 100644 testing/test_slavemanage.py delete mode 100644 xdist/nodemanage.py create mode 100644 xdist/slavemanage.py diff --git a/testing/test_gwmanage.py b/testing/test_gwmanage.py deleted file mode 100644 index d7d2208..0000000 --- a/testing/test_gwmanage.py +++ /dev/null @@ -1,121 +0,0 @@ -import py -import os -from xdist.gwmanage import GatewayManager, HostRSync -from py._test.pluginmanager import HookRelay, Registry -from py._plugin import hookspec -from xdist import newhooks -import execnet - -def pytest_funcarg__hookrecorder(request): - _pytest = request.getfuncargvalue('_pytest') - hook = request.getfuncargvalue('hook') - return _pytest.gethookrecorder(hook) - -def pytest_funcarg__hook(request): - return HookRelay([hookspec, newhooks], Registry()) - -class TestGatewayManagerPopen: - def test_popen_no_default_chdir(self, hook): - gm = GatewayManager(["popen"], hook) - assert gm.specs[0].chdir is None - - def test_default_chdir(self, hook): - l = ["ssh=noco", "socket=xyz"] - for spec in GatewayManager(l, hook).specs: - assert spec.chdir == "pyexecnetcache" - for spec in GatewayManager(l, hook, defaultchdir="abc").specs: - assert spec.chdir == "abc" - - def test_popen_makegateway_events(self, hook, hookrecorder, _pytest): - hm = GatewayManager(["popen"] * 2, hook) - hm.makegateways() - call = hookrecorder.popcall("pytest_gwmanage_newgateway") - assert call.gateway.spec == execnet.XSpec("popen") - assert call.gateway.id == "gw0" - assert call.platinfo.executable == call.gateway._rinfo().executable - call = hookrecorder.popcall("pytest_gwmanage_newgateway") - assert call.gateway.id == "gw1" - assert len(hm.group) == 2 - hm.exit() - assert not len(hm.group) - - def test_popens_rsync(self, hook, mysetup): - source = mysetup.source - hm = GatewayManager(["popen"] * 2, hook) - hm.makegateways() - assert len(hm.group) == 2 - for gw in hm.group: - class pseudoexec: - args = [] - def __init__(self, *args): - self.args.extend(args) - def waitclose(self): - pass - gw.remote_exec = pseudoexec - l = [] - hm.rsync(source, notify=lambda *args: l.append(args)) - assert not l - hm.exit() - assert not len(hm.group) - assert "sys.path.insert" in gw.remote_exec.args[0] - - def test_rsync_popen_with_path(self, hook, mysetup): - source, dest = mysetup.source, mysetup.dest - hm = GatewayManager(["popen//chdir=%s" %dest] * 1, hook) - hm.makegateways() - source.ensure("dir1", "dir2", "hello") - l = [] - hm.rsync(source, notify=lambda *args: l.append(args)) - assert len(l) == 1 - assert l[0] == ("rsyncrootready", hm.group['gw0'].spec, source) - hm.exit() - dest = dest.join(source.basename) - assert dest.join("dir1").check() - assert dest.join("dir1", "dir2").check() - assert dest.join("dir1", "dir2", 'hello').check() - - def test_rsync_same_popen_twice(self, hook, mysetup, hookrecorder): - source, dest = mysetup.source, mysetup.dest - hm = GatewayManager(["popen//chdir=%s" %dest] * 2, hook) - hm.makegateways() - source.ensure("dir1", "dir2", "hello") - hm.rsync(source) - call = hookrecorder.popcall("pytest_gwmanage_rsyncstart") - assert call.source == source - assert len(call.gateways) == 1 - assert call.gateways[0] in hm.group - call = hookrecorder.popcall("pytest_gwmanage_rsyncfinish") - -class pytest_funcarg__mysetup: - def __init__(self, request): - tmp = request.getfuncargvalue('tmpdir') - self.source = tmp.mkdir("source") - self.dest = tmp.mkdir("dest") - -class TestHRSync: - def test_hrsync_filter(self, mysetup): - source, dest = mysetup.source, mysetup.dest - source.ensure("dir", "file.txt") - source.ensure(".svn", "entries") - source.ensure(".somedotfile", "moreentries") - source.ensure("somedir", "editfile~") - syncer = HostRSync(source) - l = list(source.visit(rec=syncer.filter, - fil=syncer.filter)) - assert len(l) == 3 - basenames = [x.basename for x in l] - assert 'dir' in basenames - assert 'file.txt' in basenames - assert 'somedir' in basenames - - def test_hrsync_one_host(self, mysetup): - source, dest = mysetup.source, mysetup.dest - gw = execnet.makegateway("popen//chdir=%s" % dest) - finished = [] - rsync = HostRSync(source) - rsync.add_target_host(gw, finished=lambda: finished.append(1)) - source.join("hello.py").write("world") - rsync.send() - gw.exit() - assert dest.join(source.basename, "hello.py").check() - assert len(finished) == 1 diff --git a/testing/test_nodemanage.py b/testing/test_nodemanage.py deleted file mode 100644 index 4cbe497..0000000 --- a/testing/test_nodemanage.py +++ /dev/null @@ -1,120 +0,0 @@ -import py -from xdist.nodemanage import NodeManager - -class pytest_funcarg__mysetup: - def __init__(self, request): - basetemp = request.config.mktemp( - "mysetup-%s" % request.function.__name__, - numbered=True) - self.source = basetemp.mkdir("source") - self.dest = basetemp.mkdir("dest") - request.getfuncargvalue("_pytest") - -class TestNodeManager: - @py.test.mark.xfail - def test_rsync_roots_no_roots(self, testdir, mysetup): - mysetup.source.ensure("dir1", "file1").write("hello") - config = testdir.reparseconfig([source]) - nodemanager = NodeManager(config, ["popen//chdir=%s" % mysetup.dest]) - assert nodemanager.config.topdir == source == config.topdir - nodemanager.rsync_roots() - p, = nodemanager.gwmanager.multi_exec("import os ; channel.send(os.getcwd())").receive_each() - p = py.path.local(p) - py.builtin.print_("remote curdir", p) - assert p == mysetup.dest.join(config.topdir.basename) - assert p.join("dir1").check() - assert p.join("dir1", "file1").check() - - def test_popen_rsync_subdir(self, testdir, mysetup): - source, dest = mysetup.source, mysetup.dest - dir1 = mysetup.source.mkdir("dir1") - dir2 = dir1.mkdir("dir2") - dir2.ensure("hello") - for rsyncroot in (dir1, source): - dest.remove() - nodemanager = NodeManager(testdir.parseconfig( - "--tx", "popen//chdir=%s" % dest, - "--rsyncdir", rsyncroot, - source, - )) - assert nodemanager.config.topdir == source - nodemanager.rsync_roots() - if rsyncroot == source: - dest = dest.join("source") - assert dest.join("dir1").check() - assert dest.join("dir1", "dir2").check() - assert dest.join("dir1", "dir2", 'hello').check() - nodemanager.gwmanager.exit() - - def test_init_rsync_roots(self, testdir, mysetup): - source, dest = mysetup.source, mysetup.dest - dir2 = source.ensure("dir1", "dir2", dir=1) - source.ensure("dir1", "somefile", dir=1) - dir2.ensure("hello") - source.ensure("bogusdir", "file") - source.join("conftest.py").write(py.code.Source(""" - rsyncdirs = ['dir1/dir2'] - """)) - session = testdir.reparseconfig([source]).initsession() - nodemanager = NodeManager(session.config, ["popen//chdir=%s" % dest]) - nodemanager.rsync_roots() - assert dest.join("dir2").check() - assert not dest.join("dir1").check() - assert not dest.join("bogus").check() - - def test_rsyncignore(self, testdir, mysetup): - source, dest = mysetup.source, mysetup.dest - dir2 = source.ensure("dir1", "dir2", dir=1) - dir5 = source.ensure("dir5", "dir6", "bogus") - dirf = source.ensure("dir5", "file") - dir2.ensure("hello") - source.join("conftest.py").write(py.code.Source(""" - rsyncdirs = ['dir1', 'dir5'] - rsyncignore = ['dir1/dir2', 'dir5/dir6'] - """)) - session = testdir.reparseconfig([source]).initsession() - nodemanager = NodeManager(session.config, - ["popen//chdir=%s" % dest]) - nodemanager.rsync_roots() - assert dest.join("dir1").check() - assert not dest.join("dir1", "dir2").check() - assert dest.join("dir5","file").check() - assert not dest.join("dir6").check() - - def test_optimise_popen(self, testdir, mysetup): - source, dest = mysetup.source, mysetup.dest - specs = ["popen"] * 3 - source.join("conftest.py").write("rsyncdirs = ['a']") - source.ensure('a', dir=1) - config = testdir.reparseconfig([source]) - nodemanager = NodeManager(config, specs) - nodemanager.rsync_roots() - for gwspec in nodemanager.gwmanager.specs: - assert gwspec._samefilesystem() - assert not gwspec.chdir - - def test_setup_DEBUG(self, mysetup, testdir): - source = mysetup.source - specs = ["popen"] * 2 - source.join("conftest.py").write("rsyncdirs = ['a']") - source.ensure('a', dir=1) - config = testdir.reparseconfig([source, '--debug']) - assert config.option.debug - nodemanager = NodeManager(config, specs) - reprec = testdir.getreportrecorder(config).hookrecorder - nodemanager.setup_nodes(putevent=[].append) - for spec in nodemanager.gwmanager.specs: - l = reprec.getcalls("pytest_trace") - assert l - nodemanager.teardown_nodes() - - def test_ssh_setup_nodes(self, specssh, testdir): - testdir.makepyfile(__init__="", test_x=""" - def test_one(): - pass - """) - reprec = testdir.inline_run("-d", "--rsyncdir=%s" % testdir.tmpdir, - "--tx", specssh, testdir.tmpdir) - rep, = reprec.getreports("pytest_runtest_logreport") - assert rep.passed - diff --git a/testing/test_plugin.py b/testing/test_plugin.py index 00bfa3a..dfcb25f 100644 --- a/testing/test_plugin.py +++ b/testing/test_plugin.py @@ -1,7 +1,6 @@ import py - import execnet -from xdist.nodemanage import NodeManager +from xdist.slavemanage import NodeManager def test_dist_incompatibility_messages(testdir): result = testdir.runpytest("--pdb", "--looponfail") @@ -14,10 +13,13 @@ def test_dist_incompatibility_messages(testdir): assert "incompatible" in result.stderr.str() def test_dist_options(testdir): + from xdist.plugin import check_options config = testdir.parseconfigure("-n 2") + check_options(config) assert config.option.dist == "load" assert config.option.tx == ['popen'] * 2 config = testdir.parseconfigure("-d") + check_options(config) assert config.option.dist == "load" class TestDistOptions: @@ -40,7 +42,7 @@ class TestDistOptions: config = testdir.parseconfigure('--rsyncdir=' + str(testdir.tmpdir)) nm = NodeManager(config, specs=[execnet.XSpec("popen")]) roots = nm._getrsyncdirs() - assert len(roots) == 1 + 2 # pylib + xdist + assert len(roots) == 1 + 1 # pylib assert testdir.tmpdir in roots def test_getrsyncdirs_with_conftest(self, testdir): @@ -54,7 +56,7 @@ class TestDistOptions: testdir.tmpdir, '--rsyncdir=y', '--rsyncdir=z') nm = NodeManager(config, specs=[execnet.XSpec("popen")]) roots = nm._getrsyncdirs() - assert len(roots) == 3 + 2 # pylib + xdist + assert len(roots) == 3 + 1 # pylib assert py.path.local('y') in roots assert py.path.local('z') in roots assert testdir.tmpdir.join('x') in roots diff --git a/testing/test_remote.py b/testing/test_remote.py index 6743c75..f9ab7c4 100644 --- a/testing/test_remote.py +++ b/testing/test_remote.py @@ -1,6 +1,6 @@ import py -from xdist.remote import SlaveController -from xdist.remote import serialize_report, unserialize_report +from xdist.slavemanage import SlaveController, unserialize_report +from xdist.remote import serialize_report import execnet queue = py.builtin._tryimport("queue", "Queue") from py.builtin import print_ diff --git a/testing/test_slavemanage.py b/testing/test_slavemanage.py new file mode 100644 index 0000000..084d28e --- /dev/null +++ b/testing/test_slavemanage.py @@ -0,0 +1,241 @@ +import py +import os +import execnet +from xdist.slavemanage import GatewayManager, HostRSync +from xdist.slavemanage import NodeManager + +def pytest_funcarg__hookrecorder(request): + _pytest = request.getfuncargvalue('_pytest') + hook = request.getfuncargvalue('hook') + return _pytest.gethookrecorder(hook) + +def pytest_funcarg__hook(request): + from xdist import newhooks + from py._test.pluginmanager import HookRelay, Registry + from py._plugin import hookspec + return HookRelay([hookspec, newhooks], Registry()) + +class pytest_funcarg__mysetup: + def __init__(self, request): + basetemp = request.config.mktemp( + "mysetup-%s" % request.function.__name__, + numbered=True) + self.source = basetemp.mkdir("source") + self.dest = basetemp.mkdir("dest") + request.getfuncargvalue("_pytest") + +class TestGatewayManagerPopen: + def test_popen_no_default_chdir(self, hook): + gm = GatewayManager(["popen"], hook) + assert gm.specs[0].chdir is None + + def test_default_chdir(self, hook): + l = ["ssh=noco", "socket=xyz"] + for spec in GatewayManager(l, hook).specs: + assert spec.chdir == "pyexecnetcache" + for spec in GatewayManager(l, hook, defaultchdir="abc").specs: + assert spec.chdir == "abc" + + def test_popen_makegateway_events(self, hook, hookrecorder, _pytest): + hm = GatewayManager(["popen"] * 2, hook) + hm.makegateways() + call = hookrecorder.popcall("pytest_gwmanage_newgateway") + assert call.gateway.spec == execnet.XSpec("popen") + assert call.gateway.id == "gw0" + call = hookrecorder.popcall("pytest_gwmanage_newgateway") + assert call.gateway.id == "gw1" + assert len(hm.group) == 2 + hm.exit() + assert not len(hm.group) + + def test_popens_rsync(self, hook, mysetup): + source = mysetup.source + hm = GatewayManager(["popen"] * 2, hook) + hm.makegateways() + assert len(hm.group) == 2 + for gw in hm.group: + class pseudoexec: + args = [] + def __init__(self, *args): + self.args.extend(args) + def waitclose(self): + pass + gw.remote_exec = pseudoexec + l = [] + hm.rsync(source, notify=lambda *args: l.append(args)) + assert not l + hm.exit() + assert not len(hm.group) + assert "sys.path.insert" in gw.remote_exec.args[0] + + def test_rsync_popen_with_path(self, hook, mysetup): + source, dest = mysetup.source, mysetup.dest + hm = GatewayManager(["popen//chdir=%s" %dest] * 1, hook) + hm.makegateways() + source.ensure("dir1", "dir2", "hello") + l = [] + hm.rsync(source, notify=lambda *args: l.append(args)) + assert len(l) == 1 + assert l[0] == ("rsyncrootready", hm.group['gw0'].spec, source) + hm.exit() + dest = dest.join(source.basename) + assert dest.join("dir1").check() + assert dest.join("dir1", "dir2").check() + assert dest.join("dir1", "dir2", 'hello').check() + + def test_rsync_same_popen_twice(self, hook, mysetup, hookrecorder): + source, dest = mysetup.source, mysetup.dest + hm = GatewayManager(["popen//chdir=%s" %dest] * 2, hook) + hm.makegateways() + source.ensure("dir1", "dir2", "hello") + hm.rsync(source) + call = hookrecorder.popcall("pytest_gwmanage_rsyncstart") + assert call.source == source + assert len(call.gateways) == 1 + assert call.gateways[0] in hm.group + call = hookrecorder.popcall("pytest_gwmanage_rsyncfinish") + +class TestHRSync: + class pytest_funcarg__mysetup: + def __init__(self, request): + tmp = request.getfuncargvalue('tmpdir') + self.source = tmp.mkdir("source") + self.dest = tmp.mkdir("dest") + + def test_hrsync_filter(self, mysetup): + source, dest = mysetup.source, mysetup.dest + source.ensure("dir", "file.txt") + source.ensure(".svn", "entries") + source.ensure(".somedotfile", "moreentries") + source.ensure("somedir", "editfile~") + syncer = HostRSync(source) + l = list(source.visit(rec=syncer.filter, + fil=syncer.filter)) + assert len(l) == 3 + basenames = [x.basename for x in l] + assert 'dir' in basenames + assert 'file.txt' in basenames + assert 'somedir' in basenames + + def test_hrsync_one_host(self, mysetup): + source, dest = mysetup.source, mysetup.dest + gw = execnet.makegateway("popen//chdir=%s" % dest) + finished = [] + rsync = HostRSync(source) + rsync.add_target_host(gw, finished=lambda: finished.append(1)) + source.join("hello.py").write("world") + rsync.send() + gw.exit() + assert dest.join(source.basename, "hello.py").check() + assert len(finished) == 1 + + +class TestNodeManager: + @py.test.mark.xfail + def test_rsync_roots_no_roots(self, testdir, mysetup): + mysetup.source.ensure("dir1", "file1").write("hello") + config = testdir.reparseconfig([source]) + nodemanager = NodeManager(config, ["popen//chdir=%s" % mysetup.dest]) + #assert nodemanager.config.topdir == source == config.topdir + nodemanager.rsync_roots() + p, = nodemanager.gwmanager.multi_exec( + "import os ; channel.send(os.getcwd())").receive_each() + p = py.path.local(p) + py.builtin.print_("remote curdir", p) + assert p == mysetup.dest.join(config.topdir.basename) + assert p.join("dir1").check() + assert p.join("dir1", "file1").check() + + def test_popen_rsync_subdir(self, testdir, mysetup): + source, dest = mysetup.source, mysetup.dest + dir1 = mysetup.source.mkdir("dir1") + dir2 = dir1.mkdir("dir2") + dir2.ensure("hello") + for rsyncroot in (dir1, source): + dest.remove() + nodemanager = NodeManager(testdir.parseconfig( + "--tx", "popen//chdir=%s" % dest, + "--rsyncdir", rsyncroot, + source, + )) + #assert nodemanager.config.topdir == source + nodemanager.rsync_roots() + if rsyncroot == source: + dest = dest.join("source") + assert dest.join("dir1").check() + assert dest.join("dir1", "dir2").check() + assert dest.join("dir1", "dir2", 'hello').check() + nodemanager.gwmanager.exit() + + def test_init_rsync_roots(self, testdir, mysetup): + source, dest = mysetup.source, mysetup.dest + dir2 = source.ensure("dir1", "dir2", dir=1) + source.ensure("dir1", "somefile", dir=1) + dir2.ensure("hello") + source.ensure("bogusdir", "file") + source.join("conftest.py").write(py.code.Source(""" + rsyncdirs = ['dir1/dir2'] + """)) + config = testdir.reparseconfig([source]) + nodemanager = NodeManager(config, ["popen//chdir=%s" % dest]) + nodemanager.rsync_roots() + assert dest.join("dir2").check() + assert not dest.join("dir1").check() + assert not dest.join("bogus").check() + + def test_rsyncignore(self, testdir, mysetup): + source, dest = mysetup.source, mysetup.dest + dir2 = source.ensure("dir1", "dir2", dir=1) + dir5 = source.ensure("dir5", "dir6", "bogus") + dirf = source.ensure("dir5", "file") + dir2.ensure("hello") + source.join("conftest.py").write(py.code.Source(""" + rsyncdirs = ['dir1', 'dir5'] + rsyncignore = ['dir1/dir2', 'dir5/dir6'] + """)) + config = testdir.reparseconfig([source]) + nodemanager = NodeManager(config, ["popen//chdir=%s" % dest]) + nodemanager.rsync_roots() + assert dest.join("dir1").check() + assert not dest.join("dir1", "dir2").check() + assert dest.join("dir5","file").check() + assert not dest.join("dir6").check() + + def test_optimise_popen(self, testdir, mysetup): + source, dest = mysetup.source, mysetup.dest + specs = ["popen"] * 3 + source.join("conftest.py").write("rsyncdirs = ['a']") + source.ensure('a', dir=1) + config = testdir.reparseconfig([source]) + nodemanager = NodeManager(config, specs) + nodemanager.rsync_roots() + for gwspec in nodemanager.gwmanager.specs: + assert gwspec._samefilesystem() + assert not gwspec.chdir + + def test_setup_DEBUG(self, mysetup, testdir): + source = mysetup.source + specs = ["popen"] * 2 + source.join("conftest.py").write("rsyncdirs = ['a']") + source.ensure('a', dir=1) + config = testdir.reparseconfig([source, '--debug']) + assert config.option.debug + nodemanager = NodeManager(config, specs) + reprec = testdir.getreportrecorder(config).hookrecorder + nodemanager.setup_nodes(putevent=[].append) + for spec in nodemanager.gwmanager.specs: + l = reprec.getcalls("pytest_trace") + assert l + nodemanager.teardown_nodes() + + def test_ssh_setup_nodes(self, specssh, testdir): + testdir.makepyfile(__init__="", test_x=""" + def test_one(): + pass + """) + reprec = testdir.inline_run("-d", "--rsyncdir=%s" % testdir.tmpdir, + "--tx", specssh, testdir.tmpdir) + rep, = reprec.getreports("pytest_runtest_logreport") + assert rep.passed + + diff --git a/xdist/dsession.py b/xdist/dsession.py index 4056ee1..6bdd416 100644 --- a/xdist/dsession.py +++ b/xdist/dsession.py @@ -1,6 +1,6 @@ import py import sys -from xdist.nodemanage import NodeManager +from xdist.slavemanage import NodeManager from py._test import session queue = py.builtin._tryimport('queue', 'Queue') diff --git a/xdist/newhooks.py b/xdist/newhooks.py index 75a2df4..63ee710 100644 --- a/xdist/newhooks.py +++ b/xdist/newhooks.py @@ -1,5 +1,5 @@ -def pytest_gwmanage_newgateway(gateway, platinfo): +def pytest_gwmanage_newgateway(gateway): """ called on new raw gateway creation. """ def pytest_gwmanage_rsyncstart(source, gateways): diff --git a/xdist/nodemanage.py b/xdist/nodemanage.py deleted file mode 100644 index b7ca99a..0000000 --- a/xdist/nodemanage.py +++ /dev/null @@ -1,101 +0,0 @@ -import py -import sys, os -import xdist -from xdist.remote import SlaveController - -from xdist.gwmanage import GatewayManager -import execnet - -class NodeManager(object): - def __init__(self, config, specs=None): - self.config = config - if specs is None: - specs = self._getxspecs() - self.roots = self._getrsyncdirs() - self.gwmanager = GatewayManager(specs, config.hook) - self._nodesready = py.std.threading.Event() - - def trace(self, msg): - self.config.hook.pytest_trace(category="nodemanage", msg=msg) - - def config_getignores(self): - return self.config.getconftest_pathlist("rsyncignore") - - def rsync_roots(self): - """ make sure that all remote gateways - have the same set of roots in their - current directory. - """ - self.makegateways() - options = { - 'ignores': self.config_getignores(), - 'verbose': self.config.option.verbose, - } - if self.roots: - # send each rsync root - for root in self.roots: - self.gwmanager.rsync(root, **options) - else: - XXX # do we want to care for situations without explicit rsyncdirs? - # we transfer our topdir as the root - self.gwmanager.rsync(self.config.topdir, **options) - # and cd into it - self.gwmanager.multi_chdir(self.config.topdir.basename, inplacelocal=False) - - def makegateways(self): - # we change to the topdir sot that - # PopenGateways will have their cwd - # such that unpickling configs will - # pick it up as the right topdir - # (for other gateways this chdir is irrelevant) - self.trace("making gateways") - #old = self.config.topdir.chdir() - #try: - self.gwmanager.makegateways() - #finally: - # old.chdir() - - def setup_nodes(self, putevent): - self.rsync_roots() - self.trace("setting up nodes") - for gateway in self.gwmanager.group: - node = SlaveController(self, gateway, self.config, putevent) - gateway.node = node # to keep node alive - node.setup() - self.trace("started node %r" % node) - - def teardown_nodes(self): - self.gwmanager.exit() - - def _getxspecs(self): - config = self.config - xspeclist = [] - for xspec in config.getvalue("tx"): - i = xspec.find("*") - try: - num = int(xspec[:i]) - except ValueError: - xspeclist.append(xspec) - else: - xspeclist.extend([xspec[i+1:]] * num) - if not xspeclist: - raise config.Error( - "MISSING test execution (tx) nodes: please specify --tx") - return [execnet.XSpec(x) for x in xspeclist] - - def _getrsyncdirs(self): - config = self.config - candidates = [py._pydir] - candidates += [py.path.local(xdist.__file__).dirpath()] - candidates += config.option.rsyncdir - conftestroots = config.getconftest_pathlist("rsyncdirs") - if conftestroots: - candidates.extend(conftestroots) - roots = [] - for root in candidates: - root = py.path.local(root).realpath() - if not root.check(): - raise config.Error("rsyncdir doesn't exist: %r" %(root,)) - if root not in roots: - roots.append(root) - return roots diff --git a/xdist/remote.py b/xdist/remote.py index ff85084..e31f116 100644 --- a/xdist/remote.py +++ b/xdist/remote.py @@ -1,150 +1,13 @@ """ - Implement --dist=* testing + This module is executed in remote subprocesses and helps to + control a remote testing session and relay back information. + It assumes that 'py' is importable and does not have dependencies + on the rest of the xdist code. This means that the xdist-plugin + needs not to be installed in remote environments. """ -import py -import sys -from py._plugin import pytest_runner as runner # XXX load dynamically +import sys, os -def make_reltoroot(roots, args): - # XXX introduce/use public API for splitting py.test args - splitcode = "::" - l = [] - for arg in args: - parts = arg.split(splitcode) - fspath = py.path.local(parts[0]) - for root in roots: - x = fspath.relto(root) - if x or fspath == root: - parts[0] = root.basename + "/" + x - break - else: - raise ValueError("arg %s not relative to an rsync root" % (arg,)) - l.append(splitcode.join(parts)) - return l - -class SlaveController(object): - ENDMARK = -1 - - def __init__(self, nodemanager, gateway, config, putevent): - self.nodemanager = nodemanager - self.putevent = putevent - self.gateway = gateway - self.config = config - self.slaveinput = {'slaveid': gateway.id} - self._down = False - self.log = py.log.Producer("slavectl-%s" % gateway.id) - if not self.config.option.debug: - py.log.setconsumer(self.log._keywords, None) - - def __repr__(self): - return "<%s %s>" %(self.__class__.__name__, self.gateway.id,) - - def setup(self): - self.log("setting up slave session") - spec = self.gateway.spec - args = self.config.args - if not spec.popen or spec.chdir: - args = make_reltoroot(self.nodemanager.roots, args) - self.config.hook.pytest_configure_node(node=self) - self.channel = self.gateway.remote_exec(init_slave_session, - slaveinput=self.slaveinput, - args=args, option_dict=vars(self.config.option), - ) - if self.putevent: - self.channel.setcallback(self.process_from_remote, - endmarker=self.ENDMARK) - - def ensure_teardown(self): - if hasattr(self, 'channel'): - if not self.channel.isclosed(): - self.log("closing", self.channel) - self.channel.close() - #del self.channel - if hasattr(self, 'gateway'): - self.log("exiting", self.gateway) - self.gateway.exit() - #del self.gateway - - def send_runtest(self, nodeid): - self.sendcommand("runtests", ids=[nodeid]) - - def shutdown(self): - if not self._down and not self.channel.isclosed(): - self.sendcommand("shutdown") - - def sendcommand(self, name, **kwargs): - """ send a named parametrized command to the other side. """ - self.log("sending command %s(**%s)" % (name, kwargs)) - self.channel.send((name, kwargs)) - - def notify_inproc(self, eventname, **kwargs): - self.log("queuing %s(**%s)" % (eventname, kwargs)) - self.putevent((eventname, kwargs)) - - def process_from_remote(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_inproc("errordown", node=self, error=err) - self._down = True - return - eventname, kwargs = eventcall - if eventname in ("collectionstart"): - self.log("ignoring %s(%s)" %(eventname, kwargs)) - elif eventname == "slaveready": - self.notify_inproc(eventname, node=self, **kwargs) - elif eventname == "slavefinished": - self._down = True - self.slaveoutput = kwargs['slaveoutput'] - self.notify_inproc("slavefinished", node=self) - #elif eventname == "logstart": - # self.notify_inproc(eventname, node=self, **kwargs) - elif eventname in ("testreport", "collectreport"): - rep = unserialize_report(kwargs['data']) - self.notify_inproc(eventname, node=self, rep=rep) - elif eventname == "collectionfinish": - self.notify_inproc(eventname, node=self, ids=kwargs['ids']) - else: - raise ValueError("unknown event: %s" %(eventname,)) - 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 init_slave_session(channel, slaveinput, args, option_dict): - import py - from xdist.remote import remote_initconfig, SlaveInteractor - config = remote_initconfig(py.test.config, option_dict, args) - config.slaveinput = slaveinput - config.slaveoutput = {} - interactor = SlaveInteractor(config, channel) - config.hook.pytest_cmdline_main(config=config) - -def remote_initconfig(config, option_dict, args): - config._preparse(args) - config.option.__dict__.update(option_dict) - config.option.looponfail = False - config.option.usepdb = False - config.option.dist = "no" - config.option.distload = False - config.option.numprocesses = None - config.args = args - return config - class SlaveInteractor: def __init__(self, config, channel): self.config = config @@ -210,6 +73,7 @@ class SlaveInteractor: self.sendevent("collectreport", data=data) def serialize_report(rep): + import py d = rep.__dict__.copy() d['longrepr'] = rep.longrepr and str(rep.longrepr) or None for name in d: @@ -219,15 +83,8 @@ def serialize_report(rep): d[name] = None # for now return d -def unserialize_report(reportdict): - d = reportdict - if 'result' in d: - return runner.CollectReport(**d) - else: - return runner.TestReport(**d) - def getinfodict(): - import os, sys, platform + import platform return dict( version = sys.version, version_info = tuple(sys.version_info), @@ -236,3 +93,24 @@ def getinfodict(): executable = sys.executable, cwd = os.getcwd(), ) + +def remote_initconfig(config, option_dict, args): + config._preparse(args) + config.option.__dict__.update(option_dict) + config.option.looponfail = False + config.option.usepdb = False + config.option.dist = "no" + config.option.distload = False + config.option.numprocesses = None + config.args = args + return config + + +if __name__ == '__channelexec__': + slaveinput,args,option_dict = channel.receive() + import py + config = remote_initconfig(py.test.config, option_dict, args) + config.slaveinput = slaveinput + config.slaveoutput = {} + interactor = SlaveInteractor(config, channel) + config.hook.pytest_cmdline_main(config=config) diff --git a/xdist/slavemanage.py b/xdist/slavemanage.py new file mode 100644 index 0000000..a319e28 --- /dev/null +++ b/xdist/slavemanage.py @@ -0,0 +1,318 @@ +import py +import sys, os +import execnet +import xdist.remote + +from py._plugin import pytest_runner as runner # XXX load dynamically + +class NodeManager(object): + def __init__(self, config, specs=None): + self.config = config + if specs is None: + specs = self._getxspecs() + self.roots = self._getrsyncdirs() + self.gwmanager = GatewayManager(specs, config.hook) + self._nodesready = py.std.threading.Event() + + def trace(self, msg): + self.config.hook.pytest_trace(category="nodemanage", msg=msg) + + def config_getignores(self): + return self.config.getconftest_pathlist("rsyncignore") + + def rsync_roots(self): + """ make sure that all remote gateways + have the same set of roots in their + current directory. + """ + self.makegateways() + options = { + 'ignores': self.config_getignores(), + 'verbose': self.config.option.verbose, + } + if self.roots: + # send each rsync root + for root in self.roots: + self.gwmanager.rsync(root, **options) + else: + XXX # do we want to care for situations without explicit rsyncdirs? + # we transfer our topdir as the root + self.gwmanager.rsync(self.config.topdir, **options) + # and cd into it + self.gwmanager.multi_chdir(self.config.topdir.basename, inplacelocal=False) + + def makegateways(self): + # we change to the topdir sot that + # PopenGateways will have their cwd + # such that unpickling configs will + # pick it up as the right topdir + # (for other gateways this chdir is irrelevant) + self.trace("making gateways") + #old = self.config.topdir.chdir() + #try: + self.gwmanager.makegateways() + #finally: + # old.chdir() + + def setup_nodes(self, putevent): + self.rsync_roots() + self.trace("setting up nodes") + for gateway in self.gwmanager.group: + node = SlaveController(self, gateway, self.config, putevent) + gateway.node = node # to keep node alive + node.setup() + self.trace("started node %r" % node) + + def teardown_nodes(self): + self.gwmanager.exit() + + def _getxspecs(self): + config = self.config + xspeclist = [] + for xspec in config.getvalue("tx"): + i = xspec.find("*") + try: + num = int(xspec[:i]) + except ValueError: + xspeclist.append(xspec) + else: + xspeclist.extend([xspec[i+1:]] * num) + if not xspeclist: + raise config.Error( + "MISSING test execution (tx) nodes: please specify --tx") + return [execnet.XSpec(x) for x in xspeclist] + + def _getrsyncdirs(self): + config = self.config + candidates = [py._pydir] + candidates += config.option.rsyncdir + conftestroots = config.getconftest_pathlist("rsyncdirs") + if conftestroots: + candidates.extend(conftestroots) + roots = [] + for root in candidates: + root = py.path.local(root).realpath() + if not root.check(): + raise config.Error("rsyncdir doesn't exist: %r" %(root,)) + if root not in roots: + roots.append(root) + return roots + + +class GatewayManager: + """ + instantiating, managing and rsyncing to test hosts + """ + EXIT_TIMEOUT = 10 + def __init__(self, specs, hook, defaultchdir="pyexecnetcache"): + self.specs = [] + self.hook = hook + self.group = execnet.Group() + for spec in specs: + if not isinstance(spec, execnet.XSpec): + spec = execnet.XSpec(spec) + if not spec.chdir and not spec.popen: + spec.chdir = defaultchdir + self.specs.append(spec) + + def makegateways(self): + assert not list(self.group) + for spec in self.specs: + gw = self.group.makegateway(spec) + self.hook.pytest_gwmanage_newgateway(gateway=gw) + + def rsync(self, source, notify=None, verbose=False, ignores=None): + """ perform rsync to all remote hosts. + """ + rsync = HostRSync(source, verbose=verbose, ignores=ignores) + seen = py.builtin.set() + gateways = [] + for gateway in self.group: + spec = gateway.spec + if spec.popen and not spec.chdir: + # XXX this assumes that sources are python-packages + # and that adding the basedir does not hurt + gateway.remote_exec(""" + import sys ; sys.path.insert(0, %r) + """ % os.path.dirname(str(source))).waitclose() + continue + if spec not in seen: + def finished(): + if notify: + notify("rsyncrootready", spec, source) + rsync.add_target_host(gateway, finished=finished) + seen.add(spec) + gateways.append(gateway) + if seen: + self.hook.pytest_gwmanage_rsyncstart( + source=source, + gateways=gateways, + ) + rsync.send() + self.hook.pytest_gwmanage_rsyncfinish( + source=source, + gateways=gateways, + ) + + def exit(self): + self.group.terminate(self.EXIT_TIMEOUT) + +class HostRSync(execnet.RSync): + """ RSyncer that filters out common files + """ + def __init__(self, sourcedir, *args, **kwargs): + self._synced = {} + ignores= None + if 'ignores' in kwargs: + ignores = kwargs.pop('ignores') + self._ignores = ignores or [] + super(HostRSync, self).__init__(sourcedir=sourcedir, **kwargs) + + def filter(self, path): + path = py.path.local(path) + if not path.ext in ('.pyc', '.pyo'): + if not path.basename.endswith('~'): + if path.check(dotfile=0): + for x in self._ignores: + if path == x: + break + else: + return True + + def add_target_host(self, gateway, finished=None): + remotepath = os.path.basename(self._sourcedir) + super(HostRSync, self).add_target(gateway, remotepath, + finishedcallback=finished, + delete=True,) + + def _report_send_file(self, gateway, modified_rel_path): + if self._verbose: + path = os.path.basename(self._sourcedir) + "/" + modified_rel_path + remotepath = gateway.spec.chdir + py.builtin.print_('%s:%s <= %s' % + (gateway.spec, remotepath, path)) + + +def make_reltoroot(roots, args): + # XXX introduce/use public API for splitting py.test args + splitcode = "::" + l = [] + for arg in args: + parts = arg.split(splitcode) + fspath = py.path.local(parts[0]) + for root in roots: + x = fspath.relto(root) + if x or fspath == root: + parts[0] = root.basename + "/" + x + break + else: + raise ValueError("arg %s not relative to an rsync root" % (arg,)) + l.append(splitcode.join(parts)) + return l + +class SlaveController(object): + ENDMARK = -1 + + def __init__(self, nodemanager, gateway, config, putevent): + self.nodemanager = nodemanager + self.putevent = putevent + self.gateway = gateway + self.config = config + self.slaveinput = {'slaveid': gateway.id} + self._down = False + self.log = py.log.Producer("slavectl-%s" % gateway.id) + if not self.config.option.debug: + py.log.setconsumer(self.log._keywords, None) + + def __repr__(self): + return "<%s %s>" %(self.__class__.__name__, self.gateway.id,) + + def setup(self): + self.log("setting up slave session") + spec = self.gateway.spec + args = self.config.args + if not spec.popen or spec.chdir: + args = make_reltoroot(self.nodemanager.roots, args) + 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))) + if self.putevent: + self.channel.setcallback(self.process_from_remote, + endmarker=self.ENDMARK) + + def ensure_teardown(self): + if hasattr(self, 'channel'): + if not self.channel.isclosed(): + self.log("closing", self.channel) + self.channel.close() + #del self.channel + if hasattr(self, 'gateway'): + self.log("exiting", self.gateway) + self.gateway.exit() + #del self.gateway + + def send_runtest(self, nodeid): + self.sendcommand("runtests", ids=[nodeid]) + + def shutdown(self): + if not self._down and not self.channel.isclosed(): + self.sendcommand("shutdown") + + def sendcommand(self, name, **kwargs): + """ send a named parametrized command to the other side. """ + self.log("sending command %s(**%s)" % (name, kwargs)) + self.channel.send((name, kwargs)) + + def notify_inproc(self, eventname, **kwargs): + self.log("queuing %s(**%s)" % (eventname, kwargs)) + self.putevent((eventname, kwargs)) + + def process_from_remote(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_inproc("errordown", node=self, error=err) + self._down = True + return + eventname, kwargs = eventcall + if eventname in ("collectionstart"): + self.log("ignoring %s(%s)" %(eventname, kwargs)) + elif eventname == "slaveready": + self.notify_inproc(eventname, node=self, **kwargs) + elif eventname == "slavefinished": + self._down = True + self.slaveoutput = kwargs['slaveoutput'] + self.notify_inproc("slavefinished", node=self) + #elif eventname == "logstart": + # self.notify_inproc(eventname, node=self, **kwargs) + elif eventname in ("testreport", "collectreport"): + rep = unserialize_report(kwargs['data']) + self.notify_inproc(eventname, node=self, rep=rep) + elif eventname == "collectionfinish": + self.notify_inproc(eventname, node=self, ids=kwargs['ids']) + else: + raise ValueError("unknown event: %s" %(eventname,)) + 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 unserialize_report(reportdict): + d = reportdict + if 'result' in d: + return runner.CollectReport(**d) + else: + return runner.TestReport(**d)