Skip to content
Merged
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 @@ -61,6 +61,7 @@ private StepRuntimeSummary cloneSummary(StepRuntimeSummary summary) {
StepAction latestAction = summary.getPendingAction();
StepRuntimeSummary cloned = objectMapper.convertValue(summary, StepRuntimeSummary.class);
cloned.setPendingAction(latestAction);
summary.setPendingAction(null); // clear the pending action after passing it to step runtime
return cloned;
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -15,10 +15,12 @@
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertFalse;
import static org.junit.Assert.assertNotNull;
import static org.junit.Assert.assertNull;
import static org.junit.Assert.assertTrue;
import static org.mockito.Mockito.when;

import com.netflix.maestro.engine.MaestroEngineBaseTest;
import com.netflix.maestro.engine.db.StepAction;
import com.netflix.maestro.engine.params.DefaultParamManager;
import com.netflix.maestro.engine.params.ParamsManager;
import com.netflix.maestro.engine.steps.StepRuntime;
Expand Down Expand Up @@ -138,17 +140,22 @@ public void testStart() {
.type(StepType.NOOP)
.stepRetry(StepInstance.StepRetry.from(Defaults.DEFAULT_RETRY_POLICY))
.build();
summary.setPendingAction(StepAction.builder().build());
assertNotNull(summary.getPendingAction());
boolean ret = runtimeManager.start(workflowSummary, null, summary);
assertTrue(ret);
assertEquals(StepInstance.Status.RUNNING, summary.getRuntimeState().getStatus());
assertNotNull(summary.getRuntimeState().getExecuteTime());
assertNotNull(summary.getRuntimeState().getModifyTime());
assertEquals(1, summary.getPendingRecords().size());
assertEquals(
StepInstance.Status.NOT_CREATED, summary.getPendingRecords().get(0).getOldStatus());
assertEquals(StepInstance.Status.RUNNING, summary.getPendingRecords().get(0).getNewStatus());
StepInstance.Status.NOT_CREATED, summary.getPendingRecords().getFirst().getOldStatus());
assertEquals(
StepInstance.Status.RUNNING, summary.getPendingRecords().getFirst().getNewStatus());
assertEquals(artifact, summary.getArtifacts().get("test-artifact"));
assertTrue(summary.getTimeline().isEmpty());
// The pending action should have been cleared after passing it to step runtime
assertNull(summary.getPendingAction());
}

@Test
Expand All @@ -167,9 +174,9 @@ public void testStartFailure() {
assertNotNull(summary.getRuntimeState().getModifyTime());
assertEquals(1, summary.getPendingRecords().size());
assertEquals(
StepInstance.Status.NOT_CREATED, summary.getPendingRecords().get(0).getOldStatus());
StepInstance.Status.NOT_CREATED, summary.getPendingRecords().getFirst().getOldStatus());
assertEquals(
StepInstance.Status.USER_FAILED, summary.getPendingRecords().get(0).getNewStatus());
StepInstance.Status.USER_FAILED, summary.getPendingRecords().getFirst().getNewStatus());
assertTrue(summary.getArtifacts().isEmpty());

stepRetry.incrementByStatus(StepInstance.Status.USER_FAILED);
Expand All @@ -193,17 +200,22 @@ public void testExecute() {
.type(StepType.NOOP)
.stepRetry(StepInstance.StepRetry.from(Defaults.DEFAULT_RETRY_POLICY))
.build();
summary.setPendingAction(StepAction.builder().build());
assertNotNull(summary.getPendingAction());
boolean ret = runtimeManager.execute(workflowSummary, null, summary);
assertTrue(ret);
assertEquals(StepInstance.Status.FINISHING, summary.getRuntimeState().getStatus());
assertNotNull(summary.getRuntimeState().getFinishTime());
assertNotNull(summary.getRuntimeState().getModifyTime());
assertEquals(1, summary.getPendingRecords().size());
assertEquals(
StepInstance.Status.NOT_CREATED, summary.getPendingRecords().get(0).getOldStatus());
assertEquals(StepInstance.Status.FINISHING, summary.getPendingRecords().get(0).getNewStatus());
StepInstance.Status.NOT_CREATED, summary.getPendingRecords().getFirst().getOldStatus());
assertEquals(
StepInstance.Status.FINISHING, summary.getPendingRecords().getFirst().getNewStatus());
assertEquals(artifact, summary.getArtifacts().get("test-artifact"));
assertTrue(summary.getTimeline().isEmpty());
// The pending action should have been cleared after passing it to step runtime
assertNull(summary.getPendingAction());
}

@Test
Expand All @@ -222,9 +234,9 @@ public void testExecuteFailure() {
assertNotNull(summary.getRuntimeState().getModifyTime());
assertEquals(1, summary.getPendingRecords().size());
assertEquals(
StepInstance.Status.NOT_CREATED, summary.getPendingRecords().get(0).getOldStatus());
StepInstance.Status.NOT_CREATED, summary.getPendingRecords().getFirst().getOldStatus());
assertEquals(
StepInstance.Status.PLATFORM_FAILED, summary.getPendingRecords().get(0).getNewStatus());
StepInstance.Status.PLATFORM_FAILED, summary.getPendingRecords().getFirst().getNewStatus());
assertTrue(summary.getArtifacts().isEmpty());

stepRetry.incrementByStatus(StepInstance.Status.PLATFORM_FAILED);
Expand All @@ -248,17 +260,23 @@ public void testTerminate() {
.type(StepType.NOOP)
.stepRetry(StepInstance.StepRetry.from(Defaults.DEFAULT_RETRY_POLICY))
.build();
summary.setPendingAction(StepAction.builder().build());
assertNotNull(summary.getPendingAction());
runtimeManager.terminate(workflowSummary, summary, StepInstance.Status.STOPPED);
assertEquals(StepInstance.Status.STOPPED, summary.getRuntimeState().getStatus());
assertNotNull(summary.getRuntimeState().getEndTime());
assertNotNull(summary.getRuntimeState().getModifyTime());
assertEquals(1, summary.getPendingRecords().size());
assertEquals(
StepInstance.Status.NOT_CREATED, summary.getPendingRecords().get(0).getOldStatus());
assertEquals(StepInstance.Status.STOPPED, summary.getPendingRecords().get(0).getNewStatus());
StepInstance.Status.NOT_CREATED, summary.getPendingRecords().getFirst().getOldStatus());
assertEquals(
StepInstance.Status.STOPPED, summary.getPendingRecords().getFirst().getNewStatus());
assertEquals(artifact, summary.getArtifacts().get("test-artifact"));
assertEquals(1, summary.getTimeline().getTimelineEvents().size());
assertEquals("test termination", summary.getTimeline().getTimelineEvents().get(0).getMessage());
assertEquals(
"test termination", summary.getTimeline().getTimelineEvents().getFirst().getMessage());
// The pending action should have been cleared after passing it to step runtime
assertNull(summary.getPendingAction());
}

@Test
Expand Down
5 changes: 3 additions & 2 deletions maestro-server/src/main/resources/application-aws.yml
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,7 @@ server:
port: 8080

maestro:
queue:
queue: # check MaestroJobEvent class for queue id
properties:
1:
worker-num: 8
Expand All @@ -32,7 +32,8 @@ maestro:
scan-interval: 5000
5:
worker-num: 3
scan-interval: 10000
scan-interval: 30000
ownership-timeout: 125000 # larger than the deletion processor timeout
notifier:
type: sns
alerting:
Expand Down
2 changes: 1 addition & 1 deletion maestro-server/src/main/resources/application.yml
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@ server:
port: 8080

maestro:
queue:
queue: # check MaestroJobEvent class for queue id
properties:
1:
worker-num: 8
Expand Down