From 67e5743599a32108b164d773b2f4e730685ada21 Mon Sep 17 00:00:00 2001 From: Derek Miller Date: Mon, 22 Jun 2026 10:05:31 -0500 Subject: [PATCH] Add terminal signal dependency states (CANCELED, FAILED) behind a flag Signal dependencies only had MATCHED/PENDING/SKIPPED, so a dependency that never matched stayed PENDING forever once its owning step went terminal, and a dependency that could not be resolved had no way to surface the error. Add two terminal states to StepDependencyMatchStatus: - CANCELED: the owning step reached a terminal state while the dependency was still pending. Set by SignalHandler.onTermination (a new default no-op hook the engine calls from MaestroTask termination); MaestroSignalHandler marks still-pending dependencies CANCELED. - FAILED: resolving the dependency raised a non-retryable error, carried in the dependency details. Set via SignalDependency.markFailed(Details). The built-in matching only finds a signal or leaves it pending, so core never produces FAILED; it is the mechanism for handlers that can detect a non-retryable resolution failure. Gated by maestro.signal.terminal-states-enabled (default false): with it off, behavior is unchanged and only MATCHED/PENDING/SKIPPED are ever emitted, so existing readers are unaffected until a fleet opts in. --- .../instance/StepDependencyMatchStatus.java | 12 +++- .../models/signal/SignalDependencies.java | 38 +++++++++- .../models/signal/SignalDependenciesTest.java | 69 +++++++++++++++++++ .../engine/handlers/SignalHandler.java | 10 +++ .../maestro/engine/tasks/MaestroTask.java | 1 + .../maestro/engine/tasks/MaestroTaskTest.java | 5 +- .../config/MaestroWorkflowConfiguration.java | 11 ++- .../src/main/resources/application.yml | 3 + .../signal/handler/MaestroSignalHandler.java | 21 ++++++ .../handler/MaestroSignalHandlerTest.java | 33 ++++++++- 10 files changed, 196 insertions(+), 7 deletions(-) diff --git a/maestro-common/src/main/java/com/netflix/maestro/models/instance/StepDependencyMatchStatus.java b/maestro-common/src/main/java/com/netflix/maestro/models/instance/StepDependencyMatchStatus.java index b3776cb5..ee2c3f61 100644 --- a/maestro-common/src/main/java/com/netflix/maestro/models/instance/StepDependencyMatchStatus.java +++ b/maestro-common/src/main/java/com/netflix/maestro/models/instance/StepDependencyMatchStatus.java @@ -26,7 +26,17 @@ public enum StepDependencyMatchStatus { */ PENDING(false), /** SKIPPED status indicating that conditions are skipped by a BYPASS_STEP_DEPENDENCIES action. */ - SKIPPED(true); + SKIPPED(true), + /** + * CANCELED status indicating that the owning step reached a terminal state while the dependency + * was still pending, so it can no longer be matched. + */ + CANCELED(false), + /** + * FAILED status indicating that resolving the dependency raised a non-retryable error, carried in + * the dependency details. + */ + FAILED(false); private final boolean done; diff --git a/maestro-common/src/main/java/com/netflix/maestro/models/signal/SignalDependencies.java b/maestro-common/src/main/java/com/netflix/maestro/models/signal/SignalDependencies.java index b5a0bea6..5a3f7e90 100644 --- a/maestro-common/src/main/java/com/netflix/maestro/models/signal/SignalDependencies.java +++ b/maestro-common/src/main/java/com/netflix/maestro/models/signal/SignalDependencies.java @@ -19,6 +19,7 @@ import com.fasterxml.jackson.databind.annotation.JsonNaming; import com.netflix.maestro.annotations.Nullable; import com.netflix.maestro.models.definition.User; +import com.netflix.maestro.models.error.Details; import com.netflix.maestro.models.instance.StepDependencyMatchStatus; import com.netflix.maestro.models.timeline.TimelineEvent; import com.netflix.maestro.models.timeline.TimelineLogEvent; @@ -43,6 +44,24 @@ public boolean isSatisfied() { return dependencies.stream().allMatch(e -> e.getStatus().isDone()); } + /** + * Marks every still-pending dependency as {@link StepDependencyMatchStatus#CANCELED}, used when + * the owning step reaches a terminal state and the dependencies can no longer match. + * + * @return true if any dependency status was changed + */ + @JsonIgnore + public boolean markPendingAsCanceled() { + boolean changed = false; + for (SignalDependency dependency : dependencies) { + if (dependency.getStatus() == StepDependencyMatchStatus.PENDING) { + dependency.markCanceled(); + changed = true; + } + } + return changed; + } + /** * By passes all pending step dependencies by changing their status from PENDING to SKIPPED. It * also records the user and timestamp info in the {@link @@ -63,7 +82,7 @@ public void bypass(User user, long actionTime) { @JsonNaming(PropertyNamingStrategies.SnakeCaseStrategy.class) @JsonPropertyOrder( - value = {"name", "status", "match_params", "signal_id"}, + value = {"name", "status", "match_params", "signal_id", "details"}, alphabetic = true) @JsonInclude(JsonInclude.Include.NON_EMPTY) @Data @@ -72,6 +91,7 @@ public static class SignalDependency { private StepDependencyMatchStatus status; @Nullable private Map matchParams; @Nullable private Long signalId; + @Nullable private Details details; /** Create a new {@link SignalDependency} with PENDING match status. */ public static SignalDependency initialize(String name, Map params) { @@ -86,5 +106,21 @@ public void update(Long seqId, StepDependencyMatchStatus matchStatus) { this.signalId = seqId; this.status = matchStatus; } + + /** Marks the dependency as {@link StepDependencyMatchStatus#CANCELED}. */ + private void markCanceled() { + this.status = StepDependencyMatchStatus.CANCELED; + } + + /** + * Marks the dependency as {@link StepDependencyMatchStatus#FAILED} and records the error + * details that prevented it from being resolved. + * + * @param errorDetails the error details + */ + public void markFailed(Details errorDetails) { + this.status = StepDependencyMatchStatus.FAILED; + this.details = errorDetails; + } } } diff --git a/maestro-common/src/test/java/com/netflix/maestro/models/signal/SignalDependenciesTest.java b/maestro-common/src/test/java/com/netflix/maestro/models/signal/SignalDependenciesTest.java index 14d4df2c..667c2a1c 100644 --- a/maestro-common/src/test/java/com/netflix/maestro/models/signal/SignalDependenciesTest.java +++ b/maestro-common/src/test/java/com/netflix/maestro/models/signal/SignalDependenciesTest.java @@ -17,6 +17,7 @@ import com.netflix.maestro.MaestroBaseTest; import com.netflix.maestro.models.definition.User; +import com.netflix.maestro.models.error.Details; import com.netflix.maestro.models.instance.StepDependencyMatchStatus; import java.io.IOException; import java.util.List; @@ -84,6 +85,66 @@ public void shouldInitializeWithPendingStatus() { assertThat(signalDependencies.getInfo()).isNull(); } + @Test + public void shouldMarkPendingAsCanceled() { + boolean changed = signalDependencies.markPendingAsCanceled(); + assertThat(changed).isTrue(); + assertThat(signalDependencies.isSatisfied()).isFalse(); + assertThat(signalDependencies.getDependencies()) + .first() + .matches(i -> i.getStatus() == StepDependencyMatchStatus.CANCELED); + } + + @Test + public void shouldOnlyMarkPendingDependenciesAsCanceled() { + var matched = dependency("matched_signal", StepDependencyMatchStatus.MATCHED); + var skipped = dependency("skipped_signal", StepDependencyMatchStatus.SKIPPED); + var pending = dependency("pending_signal", StepDependencyMatchStatus.PENDING); + var dependencies = new SignalDependencies(); + dependencies.setDependencies(List.of(matched, skipped, pending)); + + boolean changed = dependencies.markPendingAsCanceled(); + + assertThat(changed).isTrue(); + assertThat(matched.getStatus()).isEqualTo(StepDependencyMatchStatus.MATCHED); + assertThat(skipped.getStatus()).isEqualTo(StepDependencyMatchStatus.SKIPPED); + assertThat(pending.getStatus()).isEqualTo(StepDependencyMatchStatus.CANCELED); + } + + @Test + public void shouldReportNoChangeWhenNoPendingDependencies() { + var matched = dependency("matched_signal", StepDependencyMatchStatus.MATCHED); + var skipped = dependency("skipped_signal", StepDependencyMatchStatus.SKIPPED); + var dependencies = new SignalDependencies(); + dependencies.setDependencies(List.of(matched, skipped)); + + assertThat(dependencies.markPendingAsCanceled()).isFalse(); + assertThat(matched.getStatus()).isEqualTo(StepDependencyMatchStatus.MATCHED); + assertThat(skipped.getStatus()).isEqualTo(StepDependencyMatchStatus.SKIPPED); + } + + @Test + public void shouldNotRemarkFailedDependencyAsCanceled() { + var dependency = signalDependencies.getDependencies().getFirst(); + dependency.markFailed(Details.create("boom")); + + boolean changed = signalDependencies.markPendingAsCanceled(); + + assertThat(changed).isFalse(); + assertThat(dependency.getStatus()).isEqualTo(StepDependencyMatchStatus.FAILED); + } + + @Test + public void shouldMarkDependencyFailedWithDetails() { + var dependency = signalDependencies.getDependencies().getFirst(); + + dependency.markFailed(Details.create("denied")); + + assertThat(dependency.getStatus()).isEqualTo(StepDependencyMatchStatus.FAILED); + assertThat(dependency.getDetails()).isNotNull(); + assertThat(dependency.getDetails().getMessage()).isEqualTo("denied"); + } + @Test public void shouldByPassStepDependencies() { User user = User.create("maestro"); @@ -95,4 +156,12 @@ public void shouldByPassStepDependencies() { .isEqualTo("Signal step dependencies have been bypassed by user [maestro]"); assertThat(signalDependencies.getInfo()).extracting("timestamp").isEqualTo(actionTime); } + + private static SignalDependencies.SignalDependency dependency( + String name, StepDependencyMatchStatus status) { + var dependency = new SignalDependencies.SignalDependency(); + dependency.setName(name); + dependency.setStatus(status); + return dependency; + } } diff --git a/maestro-engine/src/main/java/com/netflix/maestro/engine/handlers/SignalHandler.java b/maestro-engine/src/main/java/com/netflix/maestro/engine/handlers/SignalHandler.java index c9787b45..a52ca3fc 100644 --- a/maestro-engine/src/main/java/com/netflix/maestro/engine/handlers/SignalHandler.java +++ b/maestro-engine/src/main/java/com/netflix/maestro/engine/handlers/SignalHandler.java @@ -50,4 +50,14 @@ public interface SignalHandler { * @return the corresponding signal instance */ SignalInstance getSignalInstance(String signalName, long signalId); + + /** + * Called when the owning step reaches a terminal state, letting the handler finalize any + * still-pending signal dependencies (e.g. mark them canceled). Defaults to a no-op. + * + * @param workflowSummary the workflow summary + * @param stepRuntimeSummary the step runtime summary + */ + default void onTermination( + WorkflowSummary workflowSummary, StepRuntimeSummary stepRuntimeSummary) {} } diff --git a/maestro-engine/src/main/java/com/netflix/maestro/engine/tasks/MaestroTask.java b/maestro-engine/src/main/java/com/netflix/maestro/engine/tasks/MaestroTask.java index 29148c2b..f5b463ac 100644 --- a/maestro-engine/src/main/java/com/netflix/maestro/engine/tasks/MaestroTask.java +++ b/maestro-engine/src/main/java/com/netflix/maestro/engine/tasks/MaestroTask.java @@ -1131,6 +1131,7 @@ private void terminate( "step is timed out by either step or workflow timeout"); } stepRuntimeManager.terminate(workflowSummary, runtimeSummary, toStatus); + signalHandler.onTermination(workflowSummary, runtimeSummary); } } catch (RuntimeException e) { LOG.warn( diff --git a/maestro-engine/src/test/java/com/netflix/maestro/engine/tasks/MaestroTaskTest.java b/maestro-engine/src/test/java/com/netflix/maestro/engine/tasks/MaestroTaskTest.java index 4a932dfb..1dfc7c47 100644 --- a/maestro-engine/src/test/java/com/netflix/maestro/engine/tasks/MaestroTaskTest.java +++ b/maestro-engine/src/test/java/com/netflix/maestro/engine/tasks/MaestroTaskTest.java @@ -28,6 +28,7 @@ import com.netflix.maestro.engine.execution.StepRuntimeManager; import com.netflix.maestro.engine.execution.StepRuntimeSummary; import com.netflix.maestro.engine.execution.WorkflowSummary; +import com.netflix.maestro.engine.handlers.SignalHandler; import com.netflix.maestro.engine.params.OutputDataManager; import com.netflix.maestro.flow.models.Flow; import com.netflix.maestro.flow.models.Task; @@ -372,13 +373,15 @@ private StepRuntimeSummary createAndRunMaestroTask( Step stepDef, StepRuntimeSummary input, WorkflowSummary workflowSummary) { + SignalHandler signalHandler = mock(SignalHandler.class); + when(signalHandler.sendOutputSignals(any(), any())).thenReturn(true); maestroTask = new MaestroTask( stepRuntimeManager, null, paramEvaluator, MAPPER, - null, + signalHandler, mock(OutputDataManager.class), null, actionDao, diff --git a/maestro-server/src/main/java/com/netflix/maestro/server/config/MaestroWorkflowConfiguration.java b/maestro-server/src/main/java/com/netflix/maestro/server/config/MaestroWorkflowConfiguration.java index 1fbe8c73..8d6ddf36 100644 --- a/maestro-server/src/main/java/com/netflix/maestro/server/config/MaestroWorkflowConfiguration.java +++ b/maestro-server/src/main/java/com/netflix/maestro/server/config/MaestroWorkflowConfiguration.java @@ -59,6 +59,7 @@ import lombok.extern.slf4j.Slf4j; import org.apache.bval.jsr.ApacheValidationProvider; import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.beans.factory.annotation.Value; import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; @@ -221,9 +222,13 @@ public StepSyncManager stepSyncManager(MaestroStepInstanceDao instanceDao) { } @Bean - public SignalHandler signalHandler(MaestroSignalBrokerDao brokerDao) { - LOG.info("Creating maestroSignalHandler within Spring boot..."); - return new MaestroSignalHandler(brokerDao); + public SignalHandler signalHandler( + MaestroSignalBrokerDao brokerDao, + @Value("${maestro.signal.terminal-states-enabled:false}") boolean terminalStatesEnabled) { + LOG.info( + "Creating maestroSignalHandler within Spring boot (terminalStatesEnabled={})...", + terminalStatesEnabled); + return new MaestroSignalHandler(brokerDao, terminalStatesEnabled); } @Bean diff --git a/maestro-server/src/main/resources/application.yml b/maestro-server/src/main/resources/application.yml index 6ac8394a..bd5bda17 100644 --- a/maestro-server/src/main/resources/application.yml +++ b/maestro-server/src/main/resources/application.yml @@ -50,6 +50,9 @@ maestro: step-action: action-timeout: 5000 check-interval: 500 + signal: + # when true, mark still-pending signal dependencies as CANCELED once the owning step is terminal + terminal-states-enabled: false engine: configs: diff --git a/maestro-signal/src/main/java/com/netflix/maestro/signal/handler/MaestroSignalHandler.java b/maestro-signal/src/main/java/com/netflix/maestro/signal/handler/MaestroSignalHandler.java index 66dc4bc9..df89d427 100644 --- a/maestro-signal/src/main/java/com/netflix/maestro/signal/handler/MaestroSignalHandler.java +++ b/maestro-signal/src/main/java/com/netflix/maestro/signal/handler/MaestroSignalHandler.java @@ -26,6 +26,9 @@ public class MaestroSignalHandler implements SignalHandler { private final MaestroSignalBrokerDao brokerDao; + /** When enabled, marks still-pending dependencies as {@code CANCELED} on step termination. */ + private final boolean terminalStatesEnabled; + /** * Sends output signals of a step. * @@ -132,6 +135,24 @@ public boolean signalsReady( return dependencies.isSatisfied(); } + /** + * On step termination, marks any still-pending signal dependencies as {@code CANCELED} when + * terminal states are enabled, so consumers can tell a dependency never matched rather than that + * it might still match. + */ + @Override + public void onTermination( + WorkflowSummary workflowSummary, StepRuntimeSummary stepRuntimeSummary) { + if (!terminalStatesEnabled + || stepRuntimeSummary == null + || stepRuntimeSummary.getSignalDependencies() == null) { + return; + } + if (stepRuntimeSummary.getSignalDependencies().markPendingAsCanceled()) { + stepRuntimeSummary.flagToSync(); + } + } + /** params will not be null and contain at least a `name` param. */ private SignalMatchDto toSignalMatch(SignalDependencies.SignalDependency dependency) { String signalName = dependency.getName(); diff --git a/maestro-signal/src/test/java/com/netflix/maestro/signal/handler/MaestroSignalHandlerTest.java b/maestro-signal/src/test/java/com/netflix/maestro/signal/handler/MaestroSignalHandlerTest.java index f88285be..ab62ce0d 100644 --- a/maestro-signal/src/test/java/com/netflix/maestro/signal/handler/MaestroSignalHandlerTest.java +++ b/maestro-signal/src/test/java/com/netflix/maestro/signal/handler/MaestroSignalHandlerTest.java @@ -36,7 +36,7 @@ public class MaestroSignalHandlerTest extends MaestroBaseTest { @Before public void setup() { brokerDao = mock(MaestroSignalBrokerDao.class); - handler = new MaestroSignalHandler(brokerDao); + handler = new MaestroSignalHandler(brokerDao, false); } @Test @@ -120,6 +120,37 @@ public void testSignalReadyFalse() throws Exception { assertTrue(runtimeSummary.isSynced()); } + @Test + public void testOnTerminationMarksPendingAsCanceledWhenEnabled() throws Exception { + MaestroSignalHandler enabled = new MaestroSignalHandler(brokerDao, true); + StepRuntimeSummary runtimeSummary = + loadObject( + "fixtures/execution/step-runtime-summary-with-step-dependencies.json", + StepRuntimeSummary.class); + + enabled.onTermination(new WorkflowSummary(), runtimeSummary); + + assertEquals( + StepDependencyMatchStatus.CANCELED, + runtimeSummary.getSignalDependencies().getDependencies().getLast().getStatus()); + assertFalse(runtimeSummary.isSynced()); + } + + @Test + public void testOnTerminationNoopWhenDisabled() throws Exception { + StepRuntimeSummary runtimeSummary = + loadObject( + "fixtures/execution/step-runtime-summary-with-step-dependencies.json", + StepRuntimeSummary.class); + + handler.onTermination(new WorkflowSummary(), runtimeSummary); + + assertEquals( + StepDependencyMatchStatus.PENDING, + runtimeSummary.getSignalDependencies().getDependencies().getLast().getStatus()); + assertTrue(runtimeSummary.isSynced()); + } + @Test public void testGetSignalInstance() throws Exception { SignalInstance instance = new SignalInstance();