Merge "Only notify pipelines once about semaphore release"
This commit is contained in:
+10
-3
@@ -2085,7 +2085,7 @@ class Scheduler(threading.Thread):
|
||||
pipeline.manager.removeItem(item)
|
||||
return
|
||||
|
||||
def _doSemaphoreReleaseEvent(self, event, tenant):
|
||||
def _doSemaphoreReleaseEvent(self, event, tenant, notified):
|
||||
semaphore = tenant.layout.getSemaphore(
|
||||
self.abide, event.semaphore_name)
|
||||
if semaphore.global_scope:
|
||||
@@ -2094,11 +2094,14 @@ class Scheduler(threading.Thread):
|
||||
else:
|
||||
tenants = [tenant]
|
||||
for tenant in tenants:
|
||||
if tenant.name in notified:
|
||||
continue
|
||||
for pipeline_name in tenant.layout.pipelines.keys():
|
||||
event = PipelineSemaphoreReleaseEvent()
|
||||
self.pipeline_management_events[
|
||||
tenant.name][pipeline_name].put(
|
||||
event, needs_result=False)
|
||||
notified.add(tenant.name)
|
||||
|
||||
def _areAllBuildsComplete(self):
|
||||
self.log.debug("Checking if all builds are complete")
|
||||
@@ -2661,6 +2664,9 @@ class Scheduler(threading.Thread):
|
||||
" in tenant %s", tenant.name)
|
||||
|
||||
def _process_tenant_management_queue(self, tenant):
|
||||
# Set of tenant names that were notified of
|
||||
# a semaphore release.
|
||||
semaphore_notified = set()
|
||||
for event in self.management_events[tenant.name]:
|
||||
event_forwarded = False
|
||||
try:
|
||||
@@ -2671,7 +2677,8 @@ class Scheduler(threading.Thread):
|
||||
elif isinstance(event, (PromoteEvent, ChangeManagementEvent)):
|
||||
event_forwarded = self._forward_management_event(event)
|
||||
elif isinstance(event, SemaphoreReleaseEvent):
|
||||
self._doSemaphoreReleaseEvent(event, tenant)
|
||||
self._doSemaphoreReleaseEvent(
|
||||
event, tenant, semaphore_notified)
|
||||
else:
|
||||
self.log.error("Unable to handle event %s for tenant %s",
|
||||
event, tenant.name)
|
||||
@@ -2798,7 +2805,7 @@ class Scheduler(threading.Thread):
|
||||
# MODEL_API <= 32
|
||||
# Kept for backward compatibility; semaphore release events
|
||||
# are now processed in the management event queue.
|
||||
self._doSemaphoreReleaseEvent(event, pipeline.tenant)
|
||||
self._doSemaphoreReleaseEvent(event, pipeline.tenant, set())
|
||||
else:
|
||||
self.log.error("Unable to handle event %s", event)
|
||||
|
||||
|
||||
Reference in New Issue
Block a user