From 18a30fab7d0cd02f203e79d314ee966708e303fa Mon Sep 17 00:00:00 2001 From: holger krekel Date: Mon, 27 Jan 2014 11:37:33 +0100 Subject: [PATCH] fix issue419: work with collection indices instead of node ids. This reduces network message size. --- CHANGELOG | 4 ++++ testing/test_dsession.py | 14 +++++++------- testing/test_remote.py | 2 +- xdist/dsession.py | 32 +++++++++++++++++--------------- xdist/remote.py | 39 ++++++++++++++++++++++----------------- xdist/slavemanage.py | 12 +++++++----- 6 files changed, 58 insertions(+), 45 deletions(-) diff --git a/CHANGELOG b/CHANGELOG index f5ff648..dc4d27e 100644 --- a/CHANGELOG +++ b/CHANGELOG @@ -7,6 +7,10 @@ - fix pytest issue382 - produce "pytest_runtest_logstart" event again in master. Thanks Aron Curzon. +- fix pytest issue419 by sending/receiving indices into the test + collection instead of node ids (which are not neccessarily unique + for functions parametrized with duplicate values) + 1.9 ------------------------- diff --git a/testing/test_dsession.py b/testing/test_dsession.py index af1d0be..34bfa52 100644 --- a/testing/test_dsession.py +++ b/testing/test_dsession.py @@ -64,9 +64,9 @@ class TestEachScheduling: assert sched.tests_finished() assert node1.sent == ['ALL'] assert node2.sent == ['ALL'] - sched.remove_item(node1, collection[0]) + sched.remove_item(node1, 0) assert sched.tests_finished() - sched.remove_item(node2, collection[0]) + sched.remove_item(node2, 0) assert sched.tests_finished() def test_schedule_remove_node(self): @@ -105,7 +105,7 @@ class TestLoadScheduling: assert len(node1.sent) == 1 assert len(node2.sent) == 1 x = sorted(node1.sent + node2.sent) - assert x == collection + assert x == [0, 1] sched.remove_item(node1, node1.sent[0]) sched.remove_item(node2, node2.sent[0]) assert sched.tests_finished() @@ -126,14 +126,14 @@ class TestLoadScheduling: sent1 = node1.sent sent2 = node2.sent chunkitems = col[:sched.ITEM_CHUNKSIZE] - assert sent1 == chunkitems - assert sent2 == chunkitems + assert (sent1 == [0,2] and sent2 == [1,3]) or ( + sent1 == [1,3] and sent2 == [0,2]) assert sched.node2pending[node1] == sent1 assert sched.node2pending[node2] == sent2 assert len(sched.pending) == 1 for node in (node1, node2): - for i in range(sched.ITEM_CHUNKSIZE): - sched.remove_item(node, "xyz") + for i in sched.node2pending[node]: + sched.remove_item(node, i) assert not sched.pending def test_add_remove_node(self): diff --git a/testing/test_remote.py b/testing/test_remote.py index d37f15a..fe38465 100644 --- a/testing/test_remote.py +++ b/testing/test_remote.py @@ -154,7 +154,7 @@ class TestSlaveInteractor: assert ev.kwargs['topdir'] == slave.testdir.tmpdir ids = ev.kwargs['ids'] assert len(ids) == 1 - slave.sendcommand("runtests", ids=ids) + slave.sendcommand("runtests", indices=range(len(ids))) slave.sendcommand("shutdown") ev = slave.popevent("logstart") assert ev.kwargs["nodeid"].endswith("test_func") diff --git a/xdist/dsession.py b/xdist/dsession.py index 78b0ce2..b315089 100644 --- a/xdist/dsession.py +++ b/xdist/dsession.py @@ -40,15 +40,15 @@ class EachScheduling: if len(self.node2pending) >= self.numnodes: self.collection_is_completed = True - def remove_item(self, node, item): - self.node2pending[node].remove(item) + def remove_item(self, node, item_index): + self.node2pending[node].remove(item_index) 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) + crashitem = self.node2collection[node][pending.pop(0)] # XXX what about the rest of pending? return crashitem @@ -56,7 +56,7 @@ class EachScheduling: assert self.collection_is_completed for node, pending in self.node2pending.items(): node.send_runtest_all() - pending[:] = self.node2collection[node] + pending[:] = range(len(self.node2collection[node])) class LoadScheduling: LOAD_THRESHOLD_NEWITEMS = 5 @@ -94,14 +94,15 @@ class LoadScheduling: if len(self.node2collection) >= self.numnodes: self.collection_is_completed = True - def remove_item(self, node, item): + def remove_item(self, node, item_index): node_pending = self.node2pending[node] - node_pending.remove(item) + assert item_index in node_pending, (item_index, node_pending) + node_pending.remove(item_index) # pre-load items-to-test if the node may become ready if self.pending and len(node_pending) < self.LOAD_THRESHOLD_NEWITEMS: - item = self.pending.pop(0) - node_pending.append(item) - node.send_runtest(item) + item_index = self.pending.pop(0) + node_pending.append(item_index) + node.send_runtest(item_index) self.log("items waiting for node: %d" %(len(self.pending))) #self.log("node2pending: %s" %(self.node2pending,)) @@ -110,7 +111,7 @@ class LoadScheduling: if not pending: return # the node must have crashed on the item if there are pending ones - crashitem = pending.pop(0) + crashitem = self.collection[pending.pop(0)] self.pending.extend(pending) return crashitem @@ -128,17 +129,18 @@ class LoadScheduling: # all collections are the same, good. # we now create an index - self.pending = col + self.collection = col + self.pending = range(len(col)) if not col: return 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): + for i, item_index in enumerate(self.pending): nodeindex = i % num_available node, pending = available[nodeindex] - node.send_runtest(item) - pending.append(item) + node.send_runtest(item_index) + pending.append(item_index) if i >= max_one_round: break del self.pending[:i + 1] @@ -304,7 +306,7 @@ class DSession: def slave_testreport(self, node, rep): if not (rep.passed and rep.when != "call"): if rep.when in ("setup", "call"): - self.sched.remove_item(node, rep.nodeid) + self.sched.remove_item(node, rep.item_index) #self.report_line("testreport %s: %s" %(rep.id, rep.status)) rep.node = node self.config.hook.pytest_runtest_logreport(report=rep) diff --git a/xdist/remote.py b/xdist/remote.py index dc42aa0..bfce8f0 100644 --- a/xdist/remote.py +++ b/xdist/remote.py @@ -47,39 +47,44 @@ class SlaveInteractor: name, kwargs = self.channel.receive() self.log("received command %s(**%s)" % (name, kwargs)) if name == "runtests": - ids = kwargs['ids'] - for nodeid in ids: - torun.append(self._id2item[nodeid]) + torun.extend(kwargs['indices']) elif name == "runtests_all": - torun.extend(session.items) - self.log("items to run: %s" %(len(torun))) + torun.extend(range(len(session.items))) + self.log("items to run: %s" % (torun,)) while len(torun) >= 2: - item = torun.pop(0) - nextitem = torun[0] - self.config.hook.pytest_runtest_protocol(item=item, - nextitem=nextitem) + # we store item_index so that we can pick it up from the + # runtest hooks + self.run_one_test(torun) + if name == "shutdown": while torun: - self.config.hook.pytest_runtest_protocol( - item=torun.pop(0), nextitem=None) + self.run_one_test(torun) break return True + def run_one_test(self, torun): + items = self.session.items + self.item_index = torun.pop(0) + if torun: + nextitem = items[torun[0]] + else: + nextitem = None + self.config.hook.pytest_runtest_protocol( + item=items[self.item_index], + nextitem=nextitem) + def pytest_collection_finish(self, session): - self._id2item = {} - ids = [] - for item in session.items: - self._id2item[item.nodeid] = item - ids.append(item.nodeid) self.sendevent("collectionfinish", topdir=str(session.fspath), - ids=ids) + ids=[item.nodeid for item in session.items]) def pytest_runtest_logstart(self, nodeid, location): self.sendevent("logstart", nodeid=nodeid, location=location) def pytest_runtest_logreport(self, report): data = serialize_report(report) + data["item_index"] = self.item_index + assert self.session.items[self.item_index].nodeid == report.nodeid self.sendevent("testreport", data=data) def pytest_collectreport(self, report): diff --git a/xdist/slavemanage.py b/xdist/slavemanage.py index 37071d3..96f67fa 100644 --- a/xdist/slavemanage.py +++ b/xdist/slavemanage.py @@ -241,8 +241,8 @@ class SlaveController(object): self.gateway.exit() #del self.gateway - def send_runtest(self, nodeid): - self.sendcommand("runtests", ids=[nodeid]) + def send_runtest(self, index): + self.sendcommand("runtests", indices=[index]) def send_runtest_all(self): self.sendcommand("runtests_all",) @@ -292,7 +292,10 @@ class SlaveController(object): elif eventname == "logstart": self.notify_inproc(eventname, node=self, **kwargs) elif eventname in ("testreport", "collectreport", "teardownreport"): + item_index = kwargs.pop("item_index", None) rep = unserialize_report(eventname, kwargs['data']) + if item_index is not None: + rep.item_index = item_index self.notify_inproc(eventname, node=self, rep=rep) elif eventname == "collectionfinish": self.notify_inproc(eventname, node=self, ids=kwargs['ids']) @@ -307,8 +310,7 @@ class SlaveController(object): self.config.pluginmanager.notify_exception(excinfo) def unserialize_report(name, reportdict): - d = reportdict if name == "testreport": - return runner.TestReport(**d) + return runner.TestReport(**reportdict) elif name == "collectreport": - return runner.CollectReport(**d) + return runner.CollectReport(**reportdict)