remove unneccesary item2nodes data structure for load scheduling
This commit is contained in:
@@ -95,43 +95,27 @@ class LoadScheduling:
|
|||||||
self.collection_is_completed = True
|
self.collection_is_completed = True
|
||||||
|
|
||||||
def remove_item(self, node, item):
|
def remove_item(self, node, item):
|
||||||
if item not in self.item2nodes:
|
node_pending = self.node2pending[node]
|
||||||
raise AssertionError(item, self.item2nodes)
|
node_pending.remove(item)
|
||||||
nodes = self.item2nodes[item]
|
|
||||||
if node in nodes: # the node might have gone down already
|
|
||||||
nodes.remove(node)
|
|
||||||
#if not nodes:
|
|
||||||
# del self.item2nodes[item]
|
|
||||||
pending = self.node2pending[node]
|
|
||||||
pending.remove(item)
|
|
||||||
# pre-load items-to-test if the node may become ready
|
# pre-load items-to-test if the node may become ready
|
||||||
if self.pending and len(pending) < self.LOAD_THRESHOLD_NEWITEMS:
|
if self.pending and len(node_pending) < self.LOAD_THRESHOLD_NEWITEMS:
|
||||||
item = self.pending.pop(0)
|
item = self.pending.pop(0)
|
||||||
pending.append(item)
|
node_pending.append(item)
|
||||||
self.item2nodes.setdefault(item, []).append(node)
|
|
||||||
node.send_runtest(item)
|
node.send_runtest(item)
|
||||||
self.log("items waiting for node: %d" %(len(self.pending)))
|
self.log("items waiting for node: %d" %(len(self.pending)))
|
||||||
#self.log("item2pending still executing: %s" %(self.item2nodes,))
|
|
||||||
#self.log("node2pending: %s" %(self.node2pending,))
|
#self.log("node2pending: %s" %(self.node2pending,))
|
||||||
|
|
||||||
def remove_node(self, node):
|
def remove_node(self, node):
|
||||||
pending = self.node2pending.pop(node)
|
pending = self.node2pending.pop(node)
|
||||||
# KeyError if we didn't get an addnode() yet
|
|
||||||
for item in pending:
|
|
||||||
l = self.item2nodes[item]
|
|
||||||
l.remove(node)
|
|
||||||
if not l:
|
|
||||||
del self.item2nodes[item]
|
|
||||||
if not pending:
|
if not pending:
|
||||||
return
|
return
|
||||||
|
# the node must have crashed on the item if there are pending ones
|
||||||
crashitem = pending.pop(0)
|
crashitem = pending.pop(0)
|
||||||
self.pending.extend(pending)
|
self.pending.extend(pending)
|
||||||
return crashitem
|
return crashitem
|
||||||
|
|
||||||
def init_distribute(self):
|
def init_distribute(self):
|
||||||
assert self.collection_is_completed
|
assert self.collection_is_completed
|
||||||
assert not hasattr(self, 'item2nodes')
|
|
||||||
self.item2nodes = {}
|
|
||||||
# XXX allow nodes to have different collections
|
# XXX allow nodes to have different collections
|
||||||
first_node, col = list(self.node2collection.items())[0]
|
first_node, col = list(self.node2collection.items())[0]
|
||||||
for node, collection in self.node2collection.items():
|
for node, collection in self.node2collection.items():
|
||||||
@@ -154,7 +138,6 @@ class LoadScheduling:
|
|||||||
nodeindex = i % num_available
|
nodeindex = i % num_available
|
||||||
node, pending = available[nodeindex]
|
node, pending = available[nodeindex]
|
||||||
node.send_runtest(item)
|
node.send_runtest(item)
|
||||||
self.item2nodes.setdefault(item, []).append(node)
|
|
||||||
pending.append(item)
|
pending.append(item)
|
||||||
if i >= max_one_round:
|
if i >= max_one_round:
|
||||||
break
|
break
|
||||||
|
|||||||
Reference in New Issue
Block a user