improve and cleanup node-ready handling
This commit is contained in:
@@ -25,13 +25,6 @@ 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_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):
|
def test_popen_rsync_subdir(self, testdir, mysetup):
|
||||||
source, dest = mysetup.source, mysetup.dest
|
source, dest = mysetup.source, mysetup.dest
|
||||||
dir1 = mysetup.source.mkdir("dir1")
|
dir1 = mysetup.source.mkdir("dir1")
|
||||||
|
|||||||
@@ -74,6 +74,7 @@ class DSession(Session):
|
|||||||
self.terminal = config.pluginmanager.getplugin("terminalreporter")
|
self.terminal = config.pluginmanager.getplugin("terminalreporter")
|
||||||
except KeyError:
|
except KeyError:
|
||||||
self.terminal = None
|
self.terminal = None
|
||||||
|
self._nodesready = py.std.threading.Event()
|
||||||
|
|
||||||
def report_line(self, line):
|
def report_line(self, line):
|
||||||
if self.terminal:
|
if self.terminal:
|
||||||
@@ -92,7 +93,6 @@ class DSession(Session):
|
|||||||
self.sessionstarts()
|
self.sessionstarts()
|
||||||
self.setup()
|
self.setup()
|
||||||
allitems = self.collect_all_items(colitems)
|
allitems = self.collect_all_items(colitems)
|
||||||
self.nodemanager.wait_nodesready(5.0)
|
|
||||||
exitstatus = self.loop(allitems)
|
exitstatus = self.loop(allitems)
|
||||||
self.teardown()
|
self.teardown()
|
||||||
self.sessionfinishes(exitstatus=exitstatus)
|
self.sessionfinishes(exitstatus=exitstatus)
|
||||||
@@ -108,7 +108,7 @@ class DSession(Session):
|
|||||||
if loopstate.shuttingdown:
|
if loopstate.shuttingdown:
|
||||||
return self.loop_once_shutdown(loopstate)
|
return self.loop_once_shutdown(loopstate)
|
||||||
colitems = loopstate.colitems
|
colitems = loopstate.colitems
|
||||||
if loopstate.dowork and colitems:
|
if self._nodesready.isSet() and loopstate.dowork and colitems:
|
||||||
self.triggertesting(loopstate.colitems)
|
self.triggertesting(loopstate.colitems)
|
||||||
colitems[:] = []
|
colitems[:] = []
|
||||||
# we use a timeout here so that control-C gets through
|
# we use a timeout here so that control-C gets through
|
||||||
@@ -191,6 +191,9 @@ class DSession(Session):
|
|||||||
def addnode(self, node):
|
def addnode(self, node):
|
||||||
assert node not in self.node2pending
|
assert node not in self.node2pending
|
||||||
self.node2pending[node] = []
|
self.node2pending[node] = []
|
||||||
|
if (not hasattr(self, 'nodemanager') or
|
||||||
|
len(self.node2pending) == len(self.nodemanager.gwmanager.group)):
|
||||||
|
self._nodesready.set()
|
||||||
|
|
||||||
def removenode(self, node):
|
def removenode(self, node):
|
||||||
try:
|
try:
|
||||||
|
|||||||
@@ -8,6 +8,7 @@ import execnet
|
|||||||
from execnet.gateway_base import RemoteError
|
from execnet.gateway_base import RemoteError
|
||||||
|
|
||||||
class GatewayManager:
|
class GatewayManager:
|
||||||
|
EXIT_TIMEOUT = 10
|
||||||
RemoteError = RemoteError
|
RemoteError = RemoteError
|
||||||
def __init__(self, specs, hook, defaultchdir="pyexecnetcache"):
|
def __init__(self, specs, hook, defaultchdir="pyexecnetcache"):
|
||||||
self.specs = []
|
self.specs = []
|
||||||
@@ -61,7 +62,7 @@ class GatewayManager:
|
|||||||
)
|
)
|
||||||
|
|
||||||
def exit(self):
|
def exit(self):
|
||||||
self.group.terminate()
|
self.group.terminate(self.EXIT_TIMEOUT)
|
||||||
|
|
||||||
class HostRSync(execnet.RSync):
|
class HostRSync(execnet.RSync):
|
||||||
""" RSyncer that filters out common files
|
""" RSyncer that filters out common files
|
||||||
|
|||||||
@@ -12,7 +12,6 @@ class NodeManager(object):
|
|||||||
specs = self._getxspecs()
|
specs = self._getxspecs()
|
||||||
self.roots = self._getrsyncdirs()
|
self.roots = self._getrsyncdirs()
|
||||||
self.gwmanager = GatewayManager(specs, config.hook)
|
self.gwmanager = GatewayManager(specs, config.hook)
|
||||||
self.nodes = []
|
|
||||||
self._nodesready = py.std.threading.Event()
|
self._nodesready = py.std.threading.Event()
|
||||||
|
|
||||||
def trace(self, msg):
|
def trace(self, msg):
|
||||||
@@ -59,25 +58,11 @@ class NodeManager(object):
|
|||||||
self.rsync_roots()
|
self.rsync_roots()
|
||||||
self.trace("setting up nodes")
|
self.trace("setting up nodes")
|
||||||
for gateway in self.gwmanager.group:
|
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
|
gateway.node = node # to keep node alive
|
||||||
self.trace("started node %r" % node)
|
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):
|
def teardown_nodes(self):
|
||||||
# XXX do teardown nodes?
|
|
||||||
self.gwmanager.exit()
|
self.gwmanager.exit()
|
||||||
|
|
||||||
def _getxspecs(self):
|
def _getxspecs(self):
|
||||||
|
|||||||
@@ -13,12 +13,11 @@ class TXNode(object):
|
|||||||
"""
|
"""
|
||||||
ENDMARK = -1
|
ENDMARK = -1
|
||||||
|
|
||||||
def __init__(self, gateway, config, putevent, slaveready=None):
|
def __init__(self, gateway, config, putevent):
|
||||||
self.config = config
|
self.config = config
|
||||||
self.putevent = putevent
|
self.putevent = putevent
|
||||||
self.gateway = gateway
|
self.gateway = gateway
|
||||||
self.channel = install_slave(gateway, config)
|
self.channel = install_slave(gateway, config)
|
||||||
self._sendslaveready = slaveready
|
|
||||||
self.channel.setcallback(self.callback, endmarker=self.ENDMARK)
|
self.channel.setcallback(self.callback, endmarker=self.ENDMARK)
|
||||||
self._down = False
|
self._down = False
|
||||||
|
|
||||||
@@ -50,8 +49,6 @@ class TXNode(object):
|
|||||||
return
|
return
|
||||||
eventname, args, kwargs = eventcall
|
eventname, args, kwargs = eventcall
|
||||||
if eventname == "slaveready":
|
if eventname == "slaveready":
|
||||||
if self._sendslaveready:
|
|
||||||
self._sendslaveready(self)
|
|
||||||
self.notify("pytest_testnodeready", node=self)
|
self.notify("pytest_testnodeready", node=self)
|
||||||
elif eventname == "slavefinished":
|
elif eventname == "slavefinished":
|
||||||
self._down = True
|
self._down = True
|
||||||
|
|||||||
Reference in New Issue
Block a user