Added support for exchanging data between master and slave.
This commit is contained in:
@@ -117,3 +117,29 @@ class TestDistribution:
|
||||
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.
|
||||
def pytest_testnodeready(node):
|
||||
node.slavedata['data'] = 42
|
||||
|
||||
# This hook take action on slave only.
|
||||
def pytest_sessionstart(session):
|
||||
if session.__class__.__name__ == 'SlaveNode':
|
||||
assert session.slavedata['data'] == 42
|
||||
session.slavereport['result'] = 7
|
||||
|
||||
# This hook only called on master.
|
||||
def pytest_testnodedown(node, error):
|
||||
result = node.slavereport['result']
|
||||
assert result == 7
|
||||
""")
|
||||
|
||||
p1 = testdir.makepyfile("def test_func(): pass")
|
||||
result = testdir.runpytest(p1, '-d', '--tx=popen')
|
||||
result.stdout.fnmatch_lines([
|
||||
"*popen*Python*",
|
||||
"*1 passed*"
|
||||
])
|
||||
assert result.ret == 0
|
||||
|
||||
@@ -52,7 +52,8 @@ class MySetup:
|
||||
self.gateway = execnet.makegateway(self.xspec)
|
||||
self.id += 1
|
||||
self.gateway.id = str(self.id)
|
||||
self.node = TXNode(self.gateway, self.config, putevent=self.queue.put)
|
||||
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
|
||||
|
||||
|
||||
@@ -58,7 +58,7 @@ class NodeManager(object):
|
||||
self.rsync_roots()
|
||||
self.trace("setting up nodes")
|
||||
for gateway in self.gwmanager.group:
|
||||
node = TXNode(gateway, self.config, putevent)
|
||||
node = TXNode(self, gateway, self.config, putevent)
|
||||
gateway.node = node # to keep node alive
|
||||
self.trace("started node %r" % node)
|
||||
|
||||
|
||||
@@ -13,13 +13,17 @@ class TXNode(object):
|
||||
"""
|
||||
ENDMARK = -1
|
||||
|
||||
def __init__(self, gateway, config, putevent):
|
||||
def __init__(self, nodemanager, gateway, config, putevent):
|
||||
self.nodemanager = nodemanager
|
||||
self.config = config
|
||||
self.putevent = putevent
|
||||
self.gateway = gateway
|
||||
self.channel = install_slave(gateway, config)
|
||||
self.channel.setcallback(self.callback, endmarker=self.ENDMARK)
|
||||
self._down = False
|
||||
self.slavedatasent = False
|
||||
self.slavedata = {}
|
||||
self.slavereport = {}
|
||||
|
||||
def __repr__(self):
|
||||
id = self.gateway.id
|
||||
@@ -52,6 +56,7 @@ class TXNode(object):
|
||||
self.notify("pytest_testnodeready", node=self)
|
||||
elif eventname == "slavefinished":
|
||||
self._down = True
|
||||
self.slavereport = kwargs['slavereport']
|
||||
self.notify("pytest_testnodedown", error=None, node=self)
|
||||
elif eventname in ("pytest_runtest_logreport",
|
||||
"pytest__teardown_final_logerror"):
|
||||
@@ -68,16 +73,24 @@ class TXNode(object):
|
||||
self.config.pluginmanager.notify_exception(excinfo)
|
||||
|
||||
def send(self, item):
|
||||
if not self.slavedatasent:
|
||||
self.channel.send(self.slavedata)
|
||||
self.slavedatasent = True
|
||||
assert item is not None
|
||||
self.channel.send(item)
|
||||
|
||||
def sendlist(self, itemlist):
|
||||
if not self.slavedatasent:
|
||||
self.channel.send(self.slavedata)
|
||||
self.slavedatasent = True
|
||||
self.channel.send(itemlist)
|
||||
|
||||
def shutdown(self, kill=False):
|
||||
if kill:
|
||||
self.gateway.exit()
|
||||
else:
|
||||
if not self.slavedatasent:
|
||||
self.channel.send(None)
|
||||
self.channel.send(None)
|
||||
|
||||
# setting up slave code
|
||||
@@ -106,6 +119,8 @@ def install_slave(gateway, config):
|
||||
class SlaveNode(object):
|
||||
def __init__(self, channel):
|
||||
self.channel = channel
|
||||
self.slavedata = {}
|
||||
self.slavereport = {}
|
||||
|
||||
def __repr__(self):
|
||||
return "<%s channel=%s>" %(self.__class__.__name__, self.channel)
|
||||
@@ -128,6 +143,7 @@ class SlaveNode(object):
|
||||
self.config.pluginmanager.register(self)
|
||||
self.runner = self.config.pluginmanager.getplugin("pytest_runner")
|
||||
self.sendevent("slaveready")
|
||||
self.slavedata = channel.receive()
|
||||
try:
|
||||
self.config.hook.pytest_sessionstart(session=self)
|
||||
while 1:
|
||||
@@ -149,7 +165,7 @@ class SlaveNode(object):
|
||||
self.sendevent("pytest_internalerror", excrepr=er)
|
||||
raise
|
||||
else:
|
||||
self.sendevent("slavefinished")
|
||||
self.sendevent("slavefinished", slavereport=self.slavereport)
|
||||
|
||||
def run_single(self, item):
|
||||
call = self.runner.CallInfo(item._reraiseunpicklingproblem, when='setup')
|
||||
|
||||
Reference in New Issue
Block a user