Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -92,9 +92,15 @@ fun Stage.afterStages(): List<Stage> =
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()

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<DummyTask>("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"

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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)
}
}
}
Expand All @@ -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")
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -73,16 +72,13 @@ internal interface OrcaMessageHandler<M : Message> : MessageHandler<M> {
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))
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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")
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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)
}
}
}

Expand All @@ -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)
}
}
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -887,7 +887,8 @@ object CompleteStageHandlerTest : SubjectSpek<CompleteStageHandler>({

it("signals the parent stage to run") {
verify(queue).push(ContinueParentStage(
pipeline.stageByRef("1")
pipeline.stageByRef("1"),
STAGE_BEFORE
))
}
}
Expand Down Expand Up @@ -948,12 +949,17 @@ object CompleteStageHandlerTest : SubjectSpek<CompleteStageHandler>({

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
))
}
}
}
Expand Down Expand Up @@ -1061,7 +1067,8 @@ object CompleteStageHandlerTest : SubjectSpek<CompleteStageHandler>({
}

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))
}
}

Expand Down Expand Up @@ -1097,7 +1104,8 @@ object CompleteStageHandlerTest : SubjectSpek<CompleteStageHandler>({
}

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))
}
}
}
Expand Down
Loading