From c65cd3f59705281f59602698256b92b3ec3afc27 Mon Sep 17 00:00:00 2001 From: holger krekel Date: Mon, 18 Jan 2010 02:07:22 +0100 Subject: [PATCH] improve and cleanup node-ready handling --- testing/test_nodemanage.py | 7 ------- xdist/dsession.py | 7 +++++-- xdist/gwmanage.py | 3 ++- xdist/nodemanage.py | 17 +---------------- xdist/txnode.py | 5 +---- 5 files changed, 9 insertions(+), 30 deletions(-) diff --git a/testing/test_nodemanage.py b/testing/test_nodemanage.py index 4824157..bd46c9f 100644 --- a/testing/test_nodemanage.py +++ b/testing/test_nodemanage.py @@ -25,13 +25,6 @@ class TestNodeManager: assert p.join("dir1").check() assert p.join("dir1", "file1").check() - def test_popen_nodes_are_ready(self, testdir): - nodemanager = NodeManager(testdir.parseconfig( - "--tx", "3*popen")) - - nodemanager.setup_nodes([].append) - nodemanager.wait_nodesready(timeout=10.0) - def test_popen_rsync_subdir(self, testdir, mysetup): source, dest = mysetup.source, mysetup.dest dir1 = mysetup.source.mkdir("dir1") diff --git a/xdist/dsession.py b/xdist/dsession.py index f3b6bdf..4a516bf 100644 --- a/xdist/dsession.py +++ b/xdist/dsession.py @@ -74,6 +74,7 @@ class DSession(Session): self.terminal = config.pluginmanager.getplugin("terminalreporter") except KeyError: self.terminal = None + self._nodesready = py.std.threading.Event() def report_line(self, line): if self.terminal: @@ -92,7 +93,6 @@ class DSession(Session): self.sessionstarts() self.setup() allitems = self.collect_all_items(colitems) - self.nodemanager.wait_nodesready(5.0) exitstatus = self.loop(allitems) self.teardown() self.sessionfinishes(exitstatus=exitstatus) @@ -108,7 +108,7 @@ class DSession(Session): if loopstate.shuttingdown: return self.loop_once_shutdown(loopstate) colitems = loopstate.colitems - if loopstate.dowork and colitems: + if self._nodesready.isSet() and loopstate.dowork and colitems: self.triggertesting(loopstate.colitems) colitems[:] = [] # we use a timeout here so that control-C gets through @@ -191,6 +191,9 @@ class DSession(Session): def addnode(self, node): assert node not in self.node2pending self.node2pending[node] = [] + if (not hasattr(self, 'nodemanager') or + len(self.node2pending) == len(self.nodemanager.gwmanager.group)): + self._nodesready.set() def removenode(self, node): try: diff --git a/xdist/gwmanage.py b/xdist/gwmanage.py index 979a713..ec54e70 100644 --- a/xdist/gwmanage.py +++ b/xdist/gwmanage.py @@ -8,6 +8,7 @@ import execnet from execnet.gateway_base import RemoteError class GatewayManager: + EXIT_TIMEOUT = 10 RemoteError = RemoteError def __init__(self, specs, hook, defaultchdir="pyexecnetcache"): self.specs = [] @@ -61,7 +62,7 @@ class GatewayManager: ) def exit(self): - self.group.terminate() + self.group.terminate(self.EXIT_TIMEOUT) class HostRSync(execnet.RSync): """ RSyncer that filters out common files diff --git a/xdist/nodemanage.py b/xdist/nodemanage.py index d478e0a..44e4d84 100644 --- a/xdist/nodemanage.py +++ b/xdist/nodemanage.py @@ -12,7 +12,6 @@ class NodeManager(object): specs = self._getxspecs() self.roots = self._getrsyncdirs() self.gwmanager = GatewayManager(specs, config.hook) - self.nodes = [] self._nodesready = py.std.threading.Event() def trace(self, msg): @@ -59,25 +58,11 @@ class NodeManager(object): self.rsync_roots() self.trace("setting up nodes") for gateway in self.gwmanager.group: - node = TXNode(gateway, self.config, putevent, slaveready=self._slaveready) + node = TXNode(gateway, self.config, putevent) gateway.node = node # to keep node alive self.trace("started node %r" % node) - def _slaveready(self, node): - #assert node.gateway == node.gateway - #assert node.gateway.node == node - self.nodes.append(node) - self.trace("%s slave node ready %r" % (node.gateway.id, node)) - if len(self.nodes) == len(list(self.gwmanager.group)): - self._nodesready.set() - - def wait_nodesready(self, timeout=None): - self._nodesready.wait(timeout) - if not self._nodesready.isSet(): - raise IOError("nodes did not get ready for %r secs" % timeout) - def teardown_nodes(self): - # XXX do teardown nodes? self.gwmanager.exit() def _getxspecs(self): diff --git a/xdist/txnode.py b/xdist/txnode.py index 7a88f2f..85cc350 100644 --- a/xdist/txnode.py +++ b/xdist/txnode.py @@ -13,12 +13,11 @@ class TXNode(object): """ ENDMARK = -1 - def __init__(self, gateway, config, putevent, slaveready=None): + def __init__(self, gateway, config, putevent): self.config = config self.putevent = putevent self.gateway = gateway self.channel = install_slave(gateway, config) - self._sendslaveready = slaveready self.channel.setcallback(self.callback, endmarker=self.ENDMARK) self._down = False @@ -50,8 +49,6 @@ class TXNode(object): return eventname, args, kwargs = eventcall if eventname == "slaveready": - if self._sendslaveready: - self._sendslaveready(self) self.notify("pytest_testnodeready", node=self) elif eventname == "slavefinished": self._down = True