diff --git a/orca-kotlin/src/main/kotlin/com/netflix/spinnaker/orca/ext/Stage.kt b/orca-kotlin/src/main/kotlin/com/netflix/spinnaker/orca/ext/Stage.kt index bcffd65da6..39018992d7 100644 --- a/orca-kotlin/src/main/kotlin/com/netflix/spinnaker/orca/ext/Stage.kt +++ b/orca-kotlin/src/main/kotlin/com/netflix/spinnaker/orca/ext/Stage.kt @@ -92,9 +92,15 @@ fun Stage.afterStages(): List = fun Stage.allBeforeStagesComplete(): Boolean = beforeStages().all { it.status in listOf(SUCCEEDED, FAILED_CONTINUE, SKIPPED) } +fun Stage.allAfterStagesComplete(): Boolean = + afterStages().all { it.status in listOf(SUCCEEDED, FAILED_CONTINUE, SKIPPED) } + fun Stage.anyBeforeStagesFailed(): Boolean = beforeStages().any { it.status in listOf(TERMINAL, STOPPED, CANCELED) } +fun Stage.anyAfterStagesFailed(): Boolean = + afterStages().any { it.status in listOf(TERMINAL, STOPPED, CANCELED) } + fun Stage.hasTasks(): Boolean = tasks.isNotEmpty() diff --git a/orca-queue-tck/src/main/kotlin/com/netflix/spinnaker/orca/q/Stages.kt b/orca-queue-tck/src/main/kotlin/com/netflix/spinnaker/orca/q/Stages.kt index 8f5a83c1d7..93e0d1dd7a 100644 --- a/orca-queue-tck/src/main/kotlin/com/netflix/spinnaker/orca/q/Stages.kt +++ b/orca-queue-tck/src/main/kotlin/com/netflix/spinnaker/orca/q/Stages.kt @@ -116,6 +116,26 @@ val stageWithSyntheticAfter = object : StageDefinitionBuilder { } } +val stageWithParallelAfter = object : StageDefinitionBuilder { + override fun getType() = "stageWithParallelAfter" + override fun taskGraph(stage: Stage, builder: Builder) { + builder.withTask("dummy") + } + + override fun afterStages(parent: Stage, graph: StageGraphBuilder) { + graph.add { + it.type = singleTaskStage.type + it.name = "post1" + it.context = parent.context + } + graph.add { + it.type = singleTaskStage.type + it.name = "post2" + it.context = parent.context + } + } +} + val stageWithSyntheticAfterAndNoTasks = object : StageDefinitionBuilder { override fun getType() = "stageWithSyntheticAfterAndNoTasks" diff --git a/orca-queue/src/main/kotlin/com/netflix/spinnaker/orca/q/handler/ContinueParentStageHandler.kt b/orca-queue/src/main/kotlin/com/netflix/spinnaker/orca/q/handler/ContinueParentStageHandler.kt index 16db880828..e3eaa1b1da 100644 --- a/orca-queue/src/main/kotlin/com/netflix/spinnaker/orca/q/handler/ContinueParentStageHandler.kt +++ b/orca-queue/src/main/kotlin/com/netflix/spinnaker/orca/q/handler/ContinueParentStageHandler.kt @@ -17,15 +17,12 @@ package com.netflix.spinnaker.orca.q.handler import com.netflix.spinnaker.orca.ExecutionStatus.NOT_STARTED -import com.netflix.spinnaker.orca.ext.allBeforeStagesComplete -import com.netflix.spinnaker.orca.ext.anyBeforeStagesFailed -import com.netflix.spinnaker.orca.ext.firstAfterStages -import com.netflix.spinnaker.orca.ext.hasTasks +import com.netflix.spinnaker.orca.ext.* import com.netflix.spinnaker.orca.pipeline.model.Stage +import com.netflix.spinnaker.orca.pipeline.model.SyntheticStageOwner.STAGE_BEFORE import com.netflix.spinnaker.orca.pipeline.persistence.ExecutionRepository import com.netflix.spinnaker.orca.q.CompleteStage import com.netflix.spinnaker.orca.q.ContinueParentStage -import com.netflix.spinnaker.orca.q.StartStage import com.netflix.spinnaker.orca.q.StartTask import com.netflix.spinnaker.q.Queue import org.slf4j.Logger @@ -46,14 +43,23 @@ class ContinueParentStageHandler( override fun handle(message: ContinueParentStage) { message.withStage { stage -> - if (stage.allBeforeStagesComplete()) { - when { - stage.hasTasks() -> stage.runFirstTask() - else -> queue.push(CompleteStage(stage)) + if (message.phase == STAGE_BEFORE) { + if (stage.allBeforeStagesComplete()) { + when { + stage.hasTasks() -> stage.runFirstTask() + else -> queue.push(CompleteStage(stage)) + } + } else if (!stage.anyBeforeStagesFailed()) { + log.warn("Re-queuing $message as other ${message.phase} stages are still running") + queue.push(message, retryDelay) + } + } else { + if (stage.allAfterStagesComplete()) { + queue.push(CompleteStage(stage)) + } else if (!stage.anyAfterStagesFailed()) { + log.warn("Re-queuing $message as other ${message.phase} stages are still running") + queue.push(message, retryDelay) } - } else if (!stage.anyBeforeStagesFailed()) { - log.warn("Re-queuing $message as other BEFORE stages are still running") - queue.push(message, retryDelay) } } } @@ -63,18 +69,7 @@ class ContinueParentStageHandler( if (firstTask.status == NOT_STARTED) { queue.push(StartTask(this, firstTask)) } else { - log.warn("Ignoring $messageType for ${id} as tasks are already running") - } - } - - private fun Stage.runAfterStages() { - val afterStages = firstAfterStages() - if (afterStages.all { it.status == NOT_STARTED }) { - afterStages.forEach { - queue.push(StartStage(it)) - } - } else { - log.warn("Ignoring $messageType for ${id} as AFTER stages are already running") + log.warn("Ignoring $messageType for $id as tasks are already running") } } diff --git a/orca-queue/src/main/kotlin/com/netflix/spinnaker/orca/q/handler/OrcaMessageHandler.kt b/orca-queue/src/main/kotlin/com/netflix/spinnaker/orca/q/handler/OrcaMessageHandler.kt index e582273fee..a792442641 100644 --- a/orca-queue/src/main/kotlin/com/netflix/spinnaker/orca/q/handler/OrcaMessageHandler.kt +++ b/orca-queue/src/main/kotlin/com/netflix/spinnaker/orca/q/handler/OrcaMessageHandler.kt @@ -21,7 +21,6 @@ import com.netflix.spinnaker.orca.exceptions.ExceptionHandler import com.netflix.spinnaker.orca.ext.parent import com.netflix.spinnaker.orca.pipeline.model.Execution import com.netflix.spinnaker.orca.pipeline.model.Stage -import com.netflix.spinnaker.orca.pipeline.model.SyntheticStageOwner import com.netflix.spinnaker.orca.pipeline.model.Task import com.netflix.spinnaker.orca.pipeline.persistence.ExecutionNotFoundException import com.netflix.spinnaker.orca.pipeline.persistence.ExecutionRepository @@ -73,16 +72,13 @@ internal interface OrcaMessageHandler : MessageHandler { fun Stage.startNext() { execution.let { execution -> val downstreamStages = downstreamStages() + val phase = syntheticStageOwner if (downstreamStages.isNotEmpty()) { downstreamStages.forEach { queue.push(StartStage(it)) } - } else if (syntheticStageOwner == SyntheticStageOwner.STAGE_BEFORE) { - queue.push(ContinueParentStage(parent())) - } else if (syntheticStageOwner == SyntheticStageOwner.STAGE_AFTER) { - parent().let { parent -> - queue.push(CompleteStage(parent)) - } + } else if (phase != null) { + queue.push(ContinueParentStage(parent(), phase)) } else { queue.push(CompleteExecution(execution)) } diff --git a/orca-queue/src/main/kotlin/com/netflix/spinnaker/orca/q/messages.kt b/orca-queue/src/main/kotlin/com/netflix/spinnaker/orca/q/messages.kt index 41272c017a..ae4dc7c7ba 100644 --- a/orca-queue/src/main/kotlin/com/netflix/spinnaker/orca/q/messages.kt +++ b/orca-queue/src/main/kotlin/com/netflix/spinnaker/orca/q/messages.kt @@ -22,6 +22,8 @@ import com.netflix.spinnaker.orca.Task import com.netflix.spinnaker.orca.pipeline.model.Execution import com.netflix.spinnaker.orca.pipeline.model.Execution.ExecutionType import com.netflix.spinnaker.orca.pipeline.model.Stage +import com.netflix.spinnaker.orca.pipeline.model.SyntheticStageOwner +import com.netflix.spinnaker.orca.pipeline.model.SyntheticStageOwner.STAGE_BEFORE import com.netflix.spinnaker.orca.pipeline.persistence.ExecutionRepository import com.netflix.spinnaker.q.Attribute import com.netflix.spinnaker.q.Message @@ -160,13 +162,17 @@ data class ContinueParentStage( override val executionType: ExecutionType, override val executionId: String, override val application: String, - override val stageId: String + override val stageId: String, + /** + * The phase that just completed, either before or after stages. + */ + val phase: SyntheticStageOwner = STAGE_BEFORE ) : Message(), StageLevel { - constructor(source: StageLevel) : - this(source.executionType, source.executionId, source.application, source.stageId) + constructor(source: StageLevel, phase: SyntheticStageOwner) : + this(source.executionType, source.executionId, source.application, source.stageId, phase) - constructor(source: Stage) : - this(source.execution.type, source.execution.id, source.execution.application, source.id) + constructor(source: Stage, phase: SyntheticStageOwner) : + this(source.execution.type, source.execution.id, source.execution.application, source.id, phase) } @JsonTypeName("completeStage") diff --git a/orca-queue/src/test/kotlin/com/netflix/spinnaker/orca/q/MessageCompatibilityTest.kt b/orca-queue/src/test/kotlin/com/netflix/spinnaker/orca/q/MessageCompatibilityTest.kt index d5048673bd..a064cc8737 100644 --- a/orca-queue/src/test/kotlin/com/netflix/spinnaker/orca/q/MessageCompatibilityTest.kt +++ b/orca-queue/src/test/kotlin/com/netflix/spinnaker/orca/q/MessageCompatibilityTest.kt @@ -18,6 +18,8 @@ package com.netflix.spinnaker.orca.q import com.fasterxml.jackson.module.kotlin.convertValue import com.netflix.spinnaker.orca.jackson.OrcaObjectMapper +import com.netflix.spinnaker.orca.pipeline.model.SyntheticStageOwner.STAGE_AFTER +import com.netflix.spinnaker.orca.pipeline.model.SyntheticStageOwner.STAGE_BEFORE import com.netflix.spinnaker.q.Message import org.assertj.core.api.Assertions.assertThat import org.jetbrains.spek.api.Spek @@ -46,6 +48,10 @@ internal object MessageCompatibilityTest : Spek({ it("doesn't blow up") { assertThat(message).isInstanceOf(ContinueParentStage::class.java) } + + it("defaults the missing field") { + assertThat((message as ContinueParentStage).phase).isEqualTo(STAGE_BEFORE) + } } } @@ -58,6 +64,10 @@ internal object MessageCompatibilityTest : Spek({ it("doesn't blow up") { assertThat(message).isInstanceOf(ContinueParentStage::class.java) } + + it("deserializes the new field") { + assertThat((message as ContinueParentStage).phase).isEqualTo(STAGE_AFTER) + } } } } diff --git a/orca-queue/src/test/kotlin/com/netflix/spinnaker/orca/q/handler/CompleteStageHandlerTest.kt b/orca-queue/src/test/kotlin/com/netflix/spinnaker/orca/q/handler/CompleteStageHandlerTest.kt index 90cf8b788a..9448c27825 100644 --- a/orca-queue/src/test/kotlin/com/netflix/spinnaker/orca/q/handler/CompleteStageHandlerTest.kt +++ b/orca-queue/src/test/kotlin/com/netflix/spinnaker/orca/q/handler/CompleteStageHandlerTest.kt @@ -887,7 +887,8 @@ object CompleteStageHandlerTest : SubjectSpek({ it("signals the parent stage to run") { verify(queue).push(ContinueParentStage( - pipeline.stageByRef("1") + pipeline.stageByRef("1"), + STAGE_BEFORE )) } } @@ -948,12 +949,17 @@ object CompleteStageHandlerTest : SubjectSpek({ afterGroup(::resetMocks) + on("receiving the message") { subject.handle(message) } - it("signals the completion of the parent stage") { - verify(queue).push(CompleteStage(pipeline.stages.first())) + it("tells the parent stage to continue") { + verify(queue) + .push(ContinueParentStage( + pipeline.stageById(message.stageId).parent!!, + STAGE_AFTER + )) } } } @@ -1061,7 +1067,8 @@ object CompleteStageHandlerTest : SubjectSpek({ } it("signals the parent stage to try to run") { - verify(queue).push(ContinueParentStage(pipeline.stageByRef("1"))) + verify(queue) + .push(ContinueParentStage(pipeline.stageByRef("1"), STAGE_BEFORE)) } } @@ -1097,7 +1104,8 @@ object CompleteStageHandlerTest : SubjectSpek({ } it("signals the parent stage to try to run") { - verify(queue).push(ContinueParentStage(pipeline.stageByRef("1"))) + verify(queue) + .push(ContinueParentStage(pipeline.stageByRef("1"), STAGE_BEFORE)) } } } diff --git a/orca-queue/src/test/kotlin/com/netflix/spinnaker/orca/q/handler/ContinueParentStageHandlerTest.kt b/orca-queue/src/test/kotlin/com/netflix/spinnaker/orca/q/handler/ContinueParentStageHandlerTest.kt index 821b128d55..a3fe70b362 100644 --- a/orca-queue/src/test/kotlin/com/netflix/spinnaker/orca/q/handler/ContinueParentStageHandlerTest.kt +++ b/orca-queue/src/test/kotlin/com/netflix/spinnaker/orca/q/handler/ContinueParentStageHandlerTest.kt @@ -21,12 +21,17 @@ import com.netflix.spinnaker.orca.ext.beforeStages import com.netflix.spinnaker.orca.fixture.pipeline import com.netflix.spinnaker.orca.fixture.stage import com.netflix.spinnaker.orca.pipeline.model.Execution.ExecutionType.PIPELINE +import com.netflix.spinnaker.orca.pipeline.model.SyntheticStageOwner.STAGE_AFTER +import com.netflix.spinnaker.orca.pipeline.model.SyntheticStageOwner.STAGE_BEFORE import com.netflix.spinnaker.orca.pipeline.persistence.ExecutionRepository import com.netflix.spinnaker.orca.q.* import com.netflix.spinnaker.q.Queue import com.netflix.spinnaker.spek.and import com.nhaarman.mockito_kotlin.* -import org.jetbrains.spek.api.dsl.* +import org.jetbrains.spek.api.dsl.describe +import org.jetbrains.spek.api.dsl.given +import org.jetbrains.spek.api.dsl.it +import org.jetbrains.spek.api.dsl.on import org.jetbrains.spek.api.lifecycle.CachingMode import org.jetbrains.spek.subject.SubjectSpek import java.time.Duration @@ -55,7 +60,7 @@ object ContinueParentStageHandlerTest : SubjectSpek( } } - val message = ContinueParentStage(pipeline.stageByRef("1")) + val message = ContinueParentStage(pipeline.stageByRef("1"), STAGE_BEFORE) beforeGroup { pipeline.stageByRef("1<1").status = status @@ -84,7 +89,7 @@ object ContinueParentStageHandlerTest : SubjectSpek( } } - val message = ContinueParentStage(pipeline.stageByRef("1")) + val message = ContinueParentStage(pipeline.stageByRef("1"), STAGE_BEFORE) beforeGroup { pipeline.stageByRef("1<1").status = status @@ -113,7 +118,7 @@ object ContinueParentStageHandlerTest : SubjectSpek( } } - val message = ContinueParentStage(pipeline.stageByRef("1")) + val message = ContinueParentStage(pipeline.stageByRef("1"), STAGE_BEFORE) beforeGroup { pipeline.stageByRef("1").beforeStages().forEach { it.status = status } @@ -163,7 +168,7 @@ object ContinueParentStageHandlerTest : SubjectSpek( } } - val message = ContinueParentStage(pipeline.stageByRef("1")) + val message = ContinueParentStage(pipeline.stageByRef("1"), STAGE_BEFORE) beforeGroup { pipeline.stageByRef("1").beforeStages().forEach { it.status = status } @@ -182,4 +187,95 @@ object ContinueParentStageHandlerTest : SubjectSpek( } } } + + listOf(SUCCEEDED, FAILED_CONTINUE).forEach { status -> + describe("running a parent stage after its after stages complete with $status") { + given("other after stages are not yet complete") { + val pipeline = pipeline { + stage { + refId = "1" + type = stageWithParallelAfter.type + stageWithParallelAfter.buildTasks(this) + stageWithParallelAfter.buildAfterStages(this) + } + } + + val message = ContinueParentStage(pipeline.stageByRef("1"), STAGE_AFTER) + + beforeGroup { + pipeline.stageByRef("1>1").status = status + pipeline.stageByRef("1>2").status = RUNNING + whenever(repository.retrieve(PIPELINE, pipeline.id)) doReturn pipeline + } + + afterGroup(::resetMocks) + + on("receiving $message") { + subject.handle(message) + } + + it("re-queues the message for later evaluation") { + verify(queue).push(message, retryDelay) + } + } + + given("another after stage failed") { + val pipeline = pipeline { + stage { + refId = "1" + type = stageWithParallelAfter.type + stageWithParallelAfter.buildTasks(this) + stageWithParallelAfter.buildAfterStages(this) + } + } + + val message = ContinueParentStage(pipeline.stageByRef("1"), STAGE_AFTER) + + beforeGroup { + pipeline.stageByRef("1>1").status = status + pipeline.stageByRef("1>2").status = TERMINAL + whenever(repository.retrieve(PIPELINE, pipeline.id)) doReturn pipeline + } + + afterGroup(::resetMocks) + + on("receiving $message") { + subject.handle(message) + } + + it("does nothing") { + verifyZeroInteractions(queue) + } + } + + given("all after stages completed") { + val pipeline = pipeline { + stage { + refId = "1" + type = stageWithParallelAfter.type + stageWithParallelAfter.buildTasks(this) + stageWithParallelAfter.buildAfterStages(this) + } + } + + val message = ContinueParentStage(pipeline.stageByRef("1"), STAGE_AFTER) + + beforeGroup { + pipeline.stageByRef("1>1").status = status + pipeline.stageByRef("1>2").status = SUCCEEDED + whenever(repository.retrieve(PIPELINE, pipeline.id)) doReturn pipeline + } + + afterGroup(::resetMocks) + + on("receiving $message") { + subject.handle(message) + } + + it("tells the stage to complete") { + verify(queue).push(CompleteStage(pipeline.stageByRef("1"))) + } + } + } + } })