From 09b0c6bcaf27a19224c852c0b1052ae580a26217 Mon Sep 17 00:00:00 2001
From: NejlaSetkic <99647005+NejlaSetkic@users.noreply.github.com>
Date: Thu, 5 May 2022 15:38:41 +0200
Subject: [PATCH] Update abbNew.java
---
abbNew.java | 3137 +++++++++++++++++++++++++++++++--------------------
1 file changed, 1883 insertions(+), 1254 deletions(-)
diff --git a/abbNew.java b/abbNew.java
index 4a6fbfff30..e4083fd137 100644
--- a/abbNew.java
+++ b/abbNew.java
@@ -1,1383 +1,2012 @@
-package com.netflix.conductor.grpc;
-
-import com.google.protobuf.Any;
-import com.google.protobuf.Value;
-import com.netflix.conductor.common.metadata.events.EventExecution;
-import com.netflix.conductor.common.metadata.events.EventHandler;
-import com.netflix.conductor.common.metadata.tasks.PollData;
-import com.netflix.conductor.common.metadata.tasks.Task;
-import com.netflix.conductor.common.metadata.tasks.TaskDef;
-import com.netflix.conductor.common.metadata.tasks.TaskExecLog;
-import com.netflix.conductor.common.metadata.tasks.TaskResult;
-import com.netflix.conductor.common.metadata.workflow.DynamicForkJoinTask;
-import com.netflix.conductor.common.metadata.workflow.DynamicForkJoinTaskList;
+/*
+ * Copyright 2022 Netflix, Inc.
+ *
+ * Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software distributed under the License is distributed on
+ * an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations under the License.
+ */
+package com.netflix.conductor.core.execution;
+
+import java.util.*;
+import java.util.function.Predicate;
+import java.util.stream.Collectors;
+
+import org.apache.commons.lang3.StringUtils;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import org.springframework.stereotype.Component;
+
+import com.netflix.conductor.annotations.Trace;
+import com.netflix.conductor.common.metadata.tasks.*;
import com.netflix.conductor.common.metadata.workflow.RerunWorkflowRequest;
import com.netflix.conductor.common.metadata.workflow.SkipTaskRequest;
-import com.netflix.conductor.common.metadata.workflow.StartWorkflowRequest;
-import com.netflix.conductor.common.metadata.workflow.SubWorkflowParams;
import com.netflix.conductor.common.metadata.workflow.WorkflowDef;
import com.netflix.conductor.common.metadata.workflow.WorkflowTask;
-import com.netflix.conductor.common.run.TaskSummary;
import com.netflix.conductor.common.run.Workflow;
-import com.netflix.conductor.common.run.WorkflowSummary;
-import com.netflix.conductor.proto.DynamicForkJoinTaskListPb;
-import com.netflix.conductor.proto.DynamicForkJoinTaskPb;
-import com.netflix.conductor.proto.EventExecutionPb;
-import com.netflix.conductor.proto.EventHandlerPb;
-import com.netflix.conductor.proto.PollDataPb;
-import com.netflix.conductor.proto.RerunWorkflowRequestPb;
-import com.netflix.conductor.proto.SkipTaskRequestPb;
-import com.netflix.conductor.proto.StartWorkflowRequestPb;
-import com.netflix.conductor.proto.SubWorkflowParamsPb;
-import com.netflix.conductor.proto.TaskDefPb;
-import com.netflix.conductor.proto.TaskExecLogPb;
-import com.netflix.conductor.proto.TaskPb;
-import com.netflix.conductor.proto.TaskResultPb;
-import com.netflix.conductor.proto.TaskSummaryPb;
-import com.netflix.conductor.proto.WorkflowDefPb;
-import com.netflix.conductor.proto.WorkflowPb;
-import com.netflix.conductor.proto.WorkflowSummaryPb;
-import com.netflix.conductor.proto.WorkflowTaskPb;
-import java.lang.IllegalArgumentException;
-import java.lang.Object;
-import java.lang.String;
-import java.util.ArrayList;
-import java.util.HashMap;
-import java.util.HashSet;
-import java.util.List;
-import java.util.Map;
-import java.util.stream.Collectors;
-import javax.annotation.Generated;
-
-@Generated("com.netflix.conductor.annotationsprocessor.protogen")
-public abstract class AbstractProtoMapper {
- public DynamicForkJoinTaskPb.DynamicForkJoinTask toProto(DynamicForkJoinTask from) {
- DynamicForkJoinTaskPb.DynamicForkJoinTask.Builder to = DynamicForkJoinTaskPb.DynamicForkJoinTask.newBuilder();
- if (from.getTaskName() != null) {
- to.setTaskName( from.getTaskName() );
- }
- if (from.getWorkflowName() != null) {
- to.setWorkflowName( from.getWorkflowName() );
- }
- if (from.getReferenceName() != null) {
- to.setReferenceName( from.getReferenceName() );
- }
- for (Map.Entry pair : from.getInput().entrySet()) {
- to.putInput( pair.getKey(), toProto( pair.getValue() ) );
- }
- if (from.getType() != null) {
- to.setType( from.getType() );
- }
- return to.build();
+import com.netflix.conductor.common.utils.RetryUtil;
+import com.netflix.conductor.common.utils.TaskUtils;
+import com.netflix.conductor.core.WorkflowContext;
+import com.netflix.conductor.core.config.ConductorProperties;
+import com.netflix.conductor.core.dal.ExecutionDAOFacade;
+import com.netflix.conductor.core.exception.ApplicationException;
+import com.netflix.conductor.core.exception.ApplicationException.Code;
+import com.netflix.conductor.core.exception.TerminateWorkflowException;
+import com.netflix.conductor.core.execution.tasks.SystemTaskRegistry;
+import com.netflix.conductor.core.execution.tasks.Terminate;
+import com.netflix.conductor.core.execution.tasks.WorkflowSystemTask;
+import com.netflix.conductor.core.listener.WorkflowStatusListener;
+import com.netflix.conductor.core.metadata.MetadataMapperService;
+import com.netflix.conductor.core.utils.IDGenerator;
+import com.netflix.conductor.core.utils.ParametersUtils;
+import com.netflix.conductor.core.utils.QueueUtils;
+import com.netflix.conductor.dao.MetadataDAO;
+import com.netflix.conductor.dao.QueueDAO;
+import com.netflix.conductor.metrics.Monitors;
+import com.netflix.conductor.model.TaskModel;
+import com.netflix.conductor.model.WorkflowModel;
+import com.netflix.conductor.service.ExecutionLockService;
+
+import com.google.common.annotations.VisibleForTesting;
+import com.google.common.base.Preconditions;
+
+import static com.netflix.conductor.core.exception.ApplicationException.Code.*;
+import static com.netflix.conductor.core.utils.Utils.DECIDER_QUEUE;
+import static com.netflix.conductor.model.TaskModel.Status.*;
+
+/** Workflow services provider interface */
+@Trace
+@Component
+public class WorkflowExecutor {
+
+ private static final Logger LOGGER = LoggerFactory.getLogger(WorkflowExecutor.class);
+ private static final int PARENT_WF_PRIORITY = 10;
+
+ private final MetadataDAO metadataDAO;
+ private final QueueDAO queueDAO;
+ private final DeciderService deciderService;
+ private final ConductorProperties properties;
+ private final MetadataMapperService metadataMapperService;
+ private final ExecutionDAOFacade executionDAOFacade;
+ private final ParametersUtils parametersUtils;
+ private final WorkflowStatusListener workflowStatusListener;
+ private final SystemTaskRegistry systemTaskRegistry;
+
+ private long activeWorkerLastPollMs;
+ private static final String CLASS_NAME = WorkflowExecutor.class.getSimpleName();
+ private final ExecutionLockService executionLockService;
+
+ private static final Predicate UNSUCCESSFUL_TERMINAL_TASK =
+ task -> !task.getStatus().isSuccessful() && task.getStatus().isTerminal();
+
+ private static final Predicate UNSUCCESSFUL_JOIN_TASK =
+ UNSUCCESSFUL_TERMINAL_TASK.and(t -> TaskType.TASK_TYPE_JOIN.equals(t.getTaskType()));
+
+ private static final Predicate NON_TERMINAL_TASK =
+ task -> !task.getStatus().isTerminal();
+
+ private final Predicate validateLastPolledTime =
+ pollData ->
+ pollData.getLastPollTime()
+ > System.currentTimeMillis() - activeWorkerLastPollMs;
+
+ public WorkflowExecutor(
+ DeciderService deciderService,
+ MetadataDAO metadataDAO,
+ QueueDAO queueDAO,
+ MetadataMapperService metadataMapperService,
+ WorkflowStatusListener workflowStatusListener,
+ ExecutionDAOFacade executionDAOFacade,
+ ConductorProperties properties,
+ ExecutionLockService executionLockService,
+ SystemTaskRegistry systemTaskRegistry,
+ ParametersUtils parametersUtils) {
+ this.deciderService = deciderService;
+ this.metadataDAO = metadataDAO;
+ this.queueDAO = queueDAO;
+ this.properties = properties;
+ this.metadataMapperService = metadataMapperService;
+ this.executionDAOFacade = executionDAOFacade;
+ this.activeWorkerLastPollMs = properties.getActiveWorkerLastPollTimeout().toMillis();
+ this.workflowStatusListener = workflowStatusListener;
+ this.executionLockService = executionLockService;
+ this.parametersUtils = parametersUtils;
+ this.systemTaskRegistry = systemTaskRegistry;
}
- public DynamicForkJoinTask fromProto(DynamicForkJoinTaskPb.DynamicForkJoinTask from) {
- DynamicForkJoinTask to = new DynamicForkJoinTask();
- to.setTaskName( from.getTaskName() );
- to.setWorkflowName( from.getWorkflowName() );
- to.setReferenceName( from.getReferenceName() );
- Map inputMap = new HashMap();
- for (Map.Entry pair : from.getInputMap().entrySet()) {
- inputMap.put( pair.getKey(), fromProto( pair.getValue() ) );
- }
- to.setInput(inputMap);
- to.setType( from.getType() );
- return to;
+ /** @throws ApplicationException */
+ public String startWorkflow(
+ String name,
+ Integer version,
+ String correlationId,
+ Map input,
+ String externalInputPayloadStoragePath) {
+ return startWorkflow(
+ name, version, correlationId, input, externalInputPayloadStoragePath, null);
}
- public DynamicForkJoinTaskListPb.DynamicForkJoinTaskList toProto(DynamicForkJoinTaskList from) {
- DynamicForkJoinTaskListPb.DynamicForkJoinTaskList.Builder to = DynamicForkJoinTaskListPb.DynamicForkJoinTaskList.newBuilder();
- for (DynamicForkJoinTask elem : from.getDynamicTasks()) {
- to.addDynamicTasks( toProto(elem) );
- }
- return to.build();
+ /** @throws ApplicationException */
+ public String startWorkflow(
+ String name,
+ Integer version,
+ String correlationId,
+ Integer priority,
+ Map input,
+ String externalInputPayloadStoragePath) {
+ return startWorkflow(
+ name,
+ version,
+ correlationId,
+ priority,
+ input,
+ externalInputPayloadStoragePath,
+ null);
}
- public DynamicForkJoinTaskList fromProto(
- DynamicForkJoinTaskListPb.DynamicForkJoinTaskList from) {
- DynamicForkJoinTaskList to = new DynamicForkJoinTaskList();
- to.setDynamicTasks( from.getDynamicTasksList().stream().map(this::fromProto).collect(Collectors.toCollection(ArrayList::new)) );
- return to;
+ /** @throws ApplicationException */
+ public String startWorkflow(
+ String name,
+ Integer version,
+ String correlationId,
+ Map input,
+ String externalInputPayloadStoragePath,
+ String event) {
+ return startWorkflow(
+ name,
+ version,
+ input,
+ externalInputPayloadStoragePath,
+ correlationId,
+ null,
+ null,
+ event);
}
- public EventExecutionPb.EventExecution toProto(EventExecution from) {
- EventExecutionPb.EventExecution.Builder to = EventExecutionPb.EventExecution.newBuilder();
- if (from.getId() != null) {
- to.setId( from.getId() );
- }
- if (from.getMessageId() != null) {
- to.setMessageId( from.getMessageId() );
- }
- if (from.getName() != null) {
- to.setName( from.getName() );
- }
- if (from.getEvent() != null) {
- to.setEvent( from.getEvent() );
- }
- to.setCreated( from.getCreated() );
- if (from.getStatus() != null) {
- to.setStatus( toProto( from.getStatus() ) );
- }
- if (from.getAction() != null) {
- to.setAction( toProto( from.getAction() ) );
- }
- for (Map.Entry pair : from.getOutput().entrySet()) {
- to.putOutput( pair.getKey(), toProto( pair.getValue() ) );
- }
- return to.build();
+ /** @throws ApplicationException */
+ public String startWorkflow(
+ String name,
+ Integer version,
+ String correlationId,
+ Integer priority,
+ Map input,
+ String externalInputPayloadStoragePath,
+ String event) {
+ return startWorkflow(
+ name,
+ version,
+ input,
+ externalInputPayloadStoragePath,
+ correlationId,
+ priority,
+ null,
+ null,
+ event,
+ null);
}
- public EventExecution fromProto(EventExecutionPb.EventExecution from) {
- EventExecution to = new EventExecution();
- to.setId( from.getId() );
- to.setMessageId( from.getMessageId() );
- to.setName( from.getName() );
- to.setEvent( from.getEvent() );
- to.setCreated( from.getCreated() );
- to.setStatus( fromProto( from.getStatus() ) );
- to.setAction( fromProto( from.getAction() ) );
- Map outputMap = new HashMap();
- for (Map.Entry pair : from.getOutputMap().entrySet()) {
- outputMap.put( pair.getKey(), fromProto( pair.getValue() ) );
- }
- to.setOutput(outputMap);
- return to;
+ /** @throws ApplicationException */
+ public String startWorkflow(
+ String name,
+ Integer version,
+ String correlationId,
+ Map input,
+ String externalInputPayloadStoragePath,
+ String event,
+ Map taskToDomain) {
+ return startWorkflow(
+ name,
+ version,
+ correlationId,
+ 0,
+ input,
+ externalInputPayloadStoragePath,
+ event,
+ taskToDomain);
}
- public EventExecutionPb.EventExecution.Status toProto(EventExecution.Status from) {
- EventExecutionPb.EventExecution.Status to;
- switch (from) {
- case IN_PROGRESS: to = EventExecutionPb.EventExecution.Status.IN_PROGRESS; break;
- case COMPLETED: to = EventExecutionPb.EventExecution.Status.COMPLETED; break;
- case FAILED: to = EventExecutionPb.EventExecution.Status.FAILED; break;
- case SKIPPED: to = EventExecutionPb.EventExecution.Status.SKIPPED; break;
- default: throw new IllegalArgumentException("Unexpected enum constant: " + from);
- }
- return to;
+ /** @throws ApplicationException */
+ public String startWorkflow(
+ String name,
+ Integer version,
+ String correlationId,
+ Integer priority,
+ Map input,
+ String externalInputPayloadStoragePath,
+ String event,
+ Map taskToDomain) {
+ return startWorkflow(
+ name,
+ version,
+ input,
+ externalInputPayloadStoragePath,
+ correlationId,
+ priority,
+ null,
+ null,
+ event,
+ taskToDomain);
}
- public EventExecution.Status fromProto(EventExecutionPb.EventExecution.Status from) {
- EventExecution.Status to;
- switch (from) {
- case IN_PROGRESS: to = EventExecution.Status.IN_PROGRESS; break;
- case COMPLETED: to = EventExecution.Status.COMPLETED; break;
- case FAILED: to = EventExecution.Status.FAILED; break;
- case SKIPPED: to = EventExecution.Status.SKIPPED; break;
- default: throw new IllegalArgumentException("Unexpected enum constant: " + from);
- }
- return to;
+ /** @throws ApplicationException */
+ public String startWorkflow(
+ String name,
+ Integer version,
+ Map input,
+ String externalInputPayloadStoragePath,
+ String correlationId,
+ String parentWorkflowId,
+ String parentWorkflowTaskId,
+ String event) {
+ return startWorkflow(
+ name,
+ version,
+ input,
+ externalInputPayloadStoragePath,
+ correlationId,
+ parentWorkflowId,
+ parentWorkflowTaskId,
+ event,
+ null);
}
- public EventHandlerPb.EventHandler toProto(EventHandler from) {
- EventHandlerPb.EventHandler.Builder to = EventHandlerPb.EventHandler.newBuilder();
- if (from.getName() != null) {
- to.setName( from.getName() );
- }
- if (from.getEvent() != null) {
- to.setEvent( from.getEvent() );
- }
- if (from.getCondition() != null) {
- to.setCondition( from.getCondition() );
- }
- for (EventHandler.Action elem : from.getActions()) {
- to.addActions( toProto(elem) );
- }
- to.setActive( from.isActive() );
- if (from.getEvaluatorType() != null) {
- to.setEvaluatorType( from.getEvaluatorType() );
- }
- return to.build();
+ /** @throws ApplicationException */
+ public String startWorkflow(
+ WorkflowDef workflowDefinition,
+ Map workflowInput,
+ String externalInputPayloadStoragePath,
+ String correlationId,
+ String event,
+ Map taskToDomain) {
+ return startWorkflow(
+ workflowDefinition,
+ workflowInput,
+ externalInputPayloadStoragePath,
+ correlationId,
+ 0,
+ event,
+ taskToDomain);
}
- public EventHandler fromProto(EventHandlerPb.EventHandler from) {
- EventHandler to = new EventHandler();
- to.setName( from.getName() );
- to.setEvent( from.getEvent() );
- to.setCondition( from.getCondition() );
- to.setActions( from.getActionsList().stream().map(this::fromProto).collect(Collectors.toCollection(ArrayList::new)) );
- to.setActive( from.getActive() );
- to.setEvaluatorType( from.getEvaluatorType() );
- return to;
+ /** @throws ApplicationException */
+ public String startWorkflow(
+ WorkflowDef workflowDefinition,
+ Map workflowInput,
+ String externalInputPayloadStoragePath,
+ String correlationId,
+ Integer priority,
+ String event,
+ Map taskToDomain) {
+ return startWorkflow(
+ workflowDefinition,
+ workflowInput,
+ externalInputPayloadStoragePath,
+ correlationId,
+ priority,
+ null,
+ null,
+ event,
+ taskToDomain);
}
- public EventHandlerPb.EventHandler.StartWorkflow toProto(EventHandler.StartWorkflow from) {
- EventHandlerPb.EventHandler.StartWorkflow.Builder to = EventHandlerPb.EventHandler.StartWorkflow.newBuilder();
- if (from.getName() != null) {
- to.setName( from.getName() );
- }
- if (from.getVersion() != null) {
- to.setVersion( from.getVersion() );
- }
- if (from.getCorrelationId() != null) {
- to.setCorrelationId( from.getCorrelationId() );
- }
- for (Map.Entry pair : from.getInput().entrySet()) {
- to.putInput( pair.getKey(), toProto( pair.getValue() ) );
- }
- if (from.getInputMessage() != null) {
- to.setInputMessage( toProto( from.getInputMessage() ) );
- }
- to.putAllTaskToDomain( from.getTaskToDomain() );
- return to.build();
+ /** @throws ApplicationException */
+ public String startWorkflow(
+ String name,
+ Integer version,
+ Map workflowInput,
+ String externalInputPayloadStoragePath,
+ String correlationId,
+ String parentWorkflowId,
+ String parentWorkflowTaskId,
+ String event,
+ Map taskToDomain) {
+ return startWorkflow(
+ name,
+ version,
+ workflowInput,
+ externalInputPayloadStoragePath,
+ correlationId,
+ 0,
+ parentWorkflowId,
+ parentWorkflowTaskId,
+ event,
+ taskToDomain);
}
- public EventHandler.StartWorkflow fromProto(EventHandlerPb.EventHandler.StartWorkflow from) {
- EventHandler.StartWorkflow to = new EventHandler.StartWorkflow();
- to.setName( from.getName() );
- to.setVersion( from.getVersion() );
- to.setCorrelationId( from.getCorrelationId() );
- Map inputMap = new HashMap();
- for (Map.Entry pair : from.getInputMap().entrySet()) {
- inputMap.put( pair.getKey(), fromProto( pair.getValue() ) );
- }
- to.setInput(inputMap);
- if (from.hasInputMessage()) {
- to.setInputMessage( fromProto( from.getInputMessage() ) );
- }
- to.setTaskToDomain( from.getTaskToDomainMap() );
- return to;
+ /** @throws ApplicationException */
+ public String startWorkflow(
+ String name,
+ Integer version,
+ Map workflowInput,
+ String externalInputPayloadStoragePath,
+ String correlationId,
+ Integer priority,
+ String parentWorkflowId,
+ String parentWorkflowTaskId,
+ String event,
+ Map taskToDomain) {
+ WorkflowDef workflowDefinition =
+ metadataMapperService.lookupForWorkflowDefinition(name, version);
+
+ return startWorkflow(
+ workflowDefinition,
+ workflowInput,
+ externalInputPayloadStoragePath,
+ correlationId,
+ priority,
+ parentWorkflowId,
+ parentWorkflowTaskId,
+ event,
+ taskToDomain);
}
- public EventHandlerPb.EventHandler.TaskDetails toProto(EventHandler.TaskDetails from) {
- EventHandlerPb.EventHandler.TaskDetails.Builder to = EventHandlerPb.EventHandler.TaskDetails.newBuilder();
- if (from.getWorkflowId() != null) {
- to.setWorkflowId( from.getWorkflowId() );
- }
- if (from.getTaskRefName() != null) {
- to.setTaskRefName( from.getTaskRefName() );
- }
- for (Map.Entry pair : from.getOutput().entrySet()) {
- to.putOutput( pair.getKey(), toProto( pair.getValue() ) );
+ /** @throws ApplicationException if validation fails */
+ public String startWorkflow(
+ WorkflowDef workflowDefinition,
+ Map workflowInput,
+ String externalInputPayloadStoragePath,
+ String correlationId,
+ Integer priority,
+ String parentWorkflowId,
+ String parentWorkflowTaskId,
+ String event,
+ Map taskToDomain) {
+
+ workflowDefinition = metadataMapperService.populateTaskDefinitions(workflowDefinition);
+
+ // perform validations
+ validateWorkflow(workflowDefinition, workflowInput, externalInputPayloadStoragePath);
+
+ // A random UUID is assigned to the work flow instance
+ String workflowId = IDGenerator.generate();
+
+ // Persist the Workflow
+ WorkflowModel workflow = new WorkflowModel();
+ workflow.setWorkflowId(workflowId);
+ workflow.setCorrelationId(correlationId);
+ workflow.setPriority(priority == null ? 0 : priority);
+ workflow.setWorkflowDefinition(workflowDefinition);
+ workflow.setStatus(WorkflowModel.Status.RUNNING);
+ workflow.setParentWorkflowId(parentWorkflowId);
+ workflow.setParentWorkflowTaskId(parentWorkflowTaskId);
+ workflow.setOwnerApp(WorkflowContext.get().getClientApp());
+ workflow.setCreateTime(System.currentTimeMillis());
+ workflow.setUpdatedBy(null);
+ workflow.setUpdatedTime(null);
+ workflow.setEvent(event);
+ workflow.setTaskToDomain(taskToDomain);
+ workflow.setVariables(workflowDefinition.getVariables());
+
+ if (workflowInput != null && !workflowInput.isEmpty()) {
+ Map parsedInput =
+ parametersUtils.getWorkflowInput(workflowDefinition, workflowInput);
+ workflow.setInput(parsedInput);
+ } else {
+ workflow.setExternalInputPayloadStoragePath(externalInputPayloadStoragePath);
+ }
+
+ try {
+ createWorkflow(workflow);
+ // then decide to see if anything needs to be done as part of the workflow
+ decide(workflowId);
+ Monitors.recordWorkflowStartSuccess(
+ workflow.getWorkflowName(),
+ String.valueOf(workflow.getWorkflowVersion()),
+ workflow.getOwnerApp());
+ return workflowId;
+ } catch (Exception e) {
+ Monitors.recordWorkflowStartError(
+ workflowDefinition.getName(), WorkflowContext.get().getClientApp());
+ LOGGER.error("Unable to start workflow: {}", workflowDefinition.getName(), e);
+
+ // It's possible the remove workflow call hits an exception as well, in that case we
+ // want to log both errors to help diagnosis.
+ try {
+ executionDAOFacade.removeWorkflow(workflowId, false);
+ } catch (Exception rwe) {
+ LOGGER.error("Could not remove the workflowId: " + workflowId, rwe);
+ }
+ throw e;
}
- if (from.getOutputMessage() != null) {
- to.setOutputMessage( toProto( from.getOutputMessage() ) );
- }
- if (from.getTaskId() != null) {
- to.setTaskId( from.getTaskId() );
- }
- return to.build();
}
- public EventHandler.TaskDetails fromProto(EventHandlerPb.EventHandler.TaskDetails from) {
- EventHandler.TaskDetails to = new EventHandler.TaskDetails();
- to.setWorkflowId( from.getWorkflowId() );
- to.setTaskRefName( from.getTaskRefName() );
- Map outputMap = new HashMap();
- for (Map.Entry pair : from.getOutputMap().entrySet()) {
- outputMap.put( pair.getKey(), fromProto( pair.getValue() ) );
- }
- to.setOutput(outputMap);
- if (from.hasOutputMessage()) {
- to.setOutputMessage( fromProto( from.getOutputMessage() ) );
+ /*
+ * Acquire and hold the lock till the workflow creation action is completed (in primary and secondary datastores).
+ * This is to ensure that workflow creation action precedes any other action on a given workflow.
+ */
+ private void createWorkflow(WorkflowModel workflow) {
+ if (!executionLockService.acquireLock(workflow.getWorkflowId())) {
+ throw new ApplicationException(
+ BACKEND_ERROR, "Error acquiring lock when creating workflow: {}");
+ }
+ try {
+ executionDAOFacade.createWorkflow(workflow);
+ LOGGER.debug(
+ "A new instance of workflow: {} created with id: {}",
+ workflow.getWorkflowName(),
+ workflow.getWorkflowId());
+ } finally {
+ executionLockService.releaseLock(workflow.getWorkflowId());
}
- to.setTaskId( from.getTaskId() );
- return to;
}
- public EventHandlerPb.EventHandler.Action toProto(EventHandler.Action from) {
- EventHandlerPb.EventHandler.Action.Builder to = EventHandlerPb.EventHandler.Action.newBuilder();
- if (from.getAction() != null) {
- to.setAction( toProto( from.getAction() ) );
- }
- if (from.getStart_workflow() != null) {
- to.setStartWorkflow( toProto( from.getStart_workflow() ) );
- }
- if (from.getComplete_task() != null) {
- to.setCompleteTask( toProto( from.getComplete_task() ) );
- }
- if (from.getFail_task() != null) {
- to.setFailTask( toProto( from.getFail_task() ) );
+ /**
+ * Performs validations for starting a workflow
+ *
+ * @throws ApplicationException if the validation fails
+ */
+ private void validateWorkflow(
+ WorkflowDef workflowDef,
+ Map workflowInput,
+ String externalStoragePath) {
+ try {
+ // Check if the input to the workflow is not null
+ if (workflowInput == null && StringUtils.isBlank(externalStoragePath)) {
+ LOGGER.error(
+ "The input for the workflow '{}' cannot be NULL", workflowDef.getName());
+ throw new ApplicationException(
+ INVALID_INPUT, "NULL input passed when starting workflow");
+ }
+ } catch (Exception e) {
+ Monitors.recordWorkflowStartError(
+ workflowDef.getName(), WorkflowContext.get().getClientApp());
+ throw e;
}
- to.setExpandInlineJson( from.isExpandInlineJSON() );
- return to.build();
}
- public EventHandler.Action fromProto(EventHandlerPb.EventHandler.Action from) {
- EventHandler.Action to = new EventHandler.Action();
- to.setAction( fromProto( from.getAction() ) );
- if (from.hasStartWorkflow()) {
- to.setStart_workflow( fromProto( from.getStartWorkflow() ) );
- }
- if (from.hasCompleteTask()) {
- to.setComplete_task( fromProto( from.getCompleteTask() ) );
- }
- if (from.hasFailTask()) {
- to.setFail_task( fromProto( from.getFailTask() ) );
- }
- to.setExpandInlineJSON( from.getExpandInlineJson() );
- return to;
+ /**
+ * @param workflowId the id of the workflow for which task callbacks are to be reset
+ * @throws ApplicationException if the workflow is in terminal state
+ */
+ public void resetCallbacksForWorkflow(String workflowId) {
+ WorkflowModel workflow = executionDAOFacade.getWorkflowModel(workflowId, true);
+ if (workflow.getStatus().isTerminal()) {
+ throw new ApplicationException(
+ CONFLICT, "Workflow is in terminal state. Status =" + workflow.getStatus());
+ }
+
+ // Get SIMPLE tasks in SCHEDULED state that have callbackAfterSeconds > 0 and set the
+ // callbackAfterSeconds to 0
+ workflow.getTasks().stream()
+ .filter(
+ task ->
+ !systemTaskRegistry.isSystemTask(task.getTaskType())
+ && SCHEDULED == task.getStatus()
+ && task.getCallbackAfterSeconds() > 0)
+ .forEach(
+ task -> {
+ if (queueDAO.resetOffsetTime(
+ QueueUtils.getQueueName(task), task.getTaskId())) {
+ task.setCallbackAfterSeconds(0);
+ executionDAOFacade.updateTask(task);
+ }
+ });
}
- public EventHandlerPb.EventHandler.Action.Type toProto(EventHandler.Action.Type from) {
- EventHandlerPb.EventHandler.Action.Type to;
- switch (from) {
- case start_workflow: to = EventHandlerPb.EventHandler.Action.Type.START_WORKFLOW; break;
- case complete_task: to = EventHandlerPb.EventHandler.Action.Type.COMPLETE_TASK; break;
- case fail_task: to = EventHandlerPb.EventHandler.Action.Type.FAIL_TASK; break;
- default: throw new IllegalArgumentException("Unexpected enum constant: " + from);
- }
- return to;
+ public String rerun(RerunWorkflowRequest request) {
+ Preconditions.checkNotNull(
+ request.getReRunFromWorkflowId(), "reRunFromWorkflowId is missing");
+ if (!rerunWF(
+ request.getReRunFromWorkflowId(),
+ request.getReRunFromTaskId(),
+ request.getTaskInput(),
+ request.getWorkflowInput(),
+ request.getCorrelationId())) {
+ throw new ApplicationException(
+ INVALID_INPUT, "Task " + request.getReRunFromTaskId() + " not found");
+ }
+ return request.getReRunFromWorkflowId();
}
- public EventHandler.Action.Type fromProto(EventHandlerPb.EventHandler.Action.Type from) {
- EventHandler.Action.Type to;
- switch (from) {
- case START_WORKFLOW: to = EventHandler.Action.Type.start_workflow; break;
- case COMPLETE_TASK: to = EventHandler.Action.Type.complete_task; break;
- case FAIL_TASK: to = EventHandler.Action.Type.fail_task; break;
- default: throw new IllegalArgumentException("Unexpected enum constant: " + from);
- }
- return to;
+ /**
+ * @param workflowId the id of the workflow to be restarted
+ * @param useLatestDefinitions if true, use the latest workflow and task definitions upon
+ * restart
+ * @throws ApplicationException in the following cases:
+ *
+ * - Workflow is not in a terminal state
+ *
- Workflow definition is not found
+ *
- Workflow is deemed non-restartable as per workflow definition
+ *
+ */
+ public void restart(String workflowId, boolean useLatestDefinitions) {
+ WorkflowModel workflow = executionDAOFacade.getWorkflowModel(workflowId, true);
+ if (!workflow.getStatus().isTerminal()) {
+ String errorMsg =
+ String.format(
+ "Workflow: %s is not in terminal state, unable to restart.", workflow);
+ LOGGER.error(errorMsg);
+ throw new ApplicationException(CONFLICT, errorMsg);
+ }
+
+ WorkflowDef workflowDef;
+ if (useLatestDefinitions) {
+ workflowDef =
+ metadataDAO
+ .getLatestWorkflowDef(workflow.getWorkflowName())
+ .orElseThrow(
+ () ->
+ new ApplicationException(
+ NOT_FOUND,
+ String.format(
+ "Unable to find latest definition for %s",
+ workflowId)));
+ workflow.setWorkflowDefinition(workflowDef);
+ } else {
+ workflowDef =
+ Optional.ofNullable(workflow.getWorkflowDefinition())
+ .orElseGet(
+ () ->
+ metadataDAO
+ .getWorkflowDef(
+ workflow.getWorkflowName(),
+ workflow.getWorkflowVersion())
+ .orElseThrow(
+ () ->
+ new ApplicationException(
+ NOT_FOUND,
+ String.format(
+ "Unable to find definition for %s",
+ workflowId))));
+ }
+
+ if (!workflowDef.isRestartable()
+ && workflow.getStatus()
+ .equals(
+ WorkflowModel.Status
+ .COMPLETED)) { // Can only restart non-completed workflows
+ // when the configuration is set to false
+ throw new ApplicationException(
+ CONFLICT, String.format("Workflow: %s is non-restartable", workflow));
+ }
+
+ // Reset the workflow in the primary datastore and remove from indexer; then re-create it
+ executionDAOFacade.resetWorkflow(workflowId);
+
+ workflow.getTasks().clear();
+ workflow.setReasonForIncompletion(null);
+ workflow.setFailedTaskId(null);
+ workflow.setCreateTime(System.currentTimeMillis());
+ workflow.setEndTime(0);
+ workflow.setLastRetriedTime(0);
+ // Change the status to running
+ workflow.setStatus(WorkflowModel.Status.RUNNING);
+ workflow.setOutput(null);
+ workflow.setExternalOutputPayloadStoragePath(null);
+
+ try {
+ executionDAOFacade.createWorkflow(workflow);
+ } catch (Exception e) {
+ Monitors.recordWorkflowStartError(
+ workflowDef.getName(), WorkflowContext.get().getClientApp());
+ LOGGER.error("Unable to restart workflow: {}", workflowDef.getName(), e);
+ terminateWorkflow(workflowId, "Error when restarting the workflow");
+ throw e;
+ }
+
+ decide(workflowId);
+
+ updateAndPushParents(workflow, "restarted");
}
- public PollDataPb.PollData toProto(PollData from) {
- PollDataPb.PollData.Builder to = PollDataPb.PollData.newBuilder();
- if (from.getQueueName() != null) {
- to.setQueueName( from.getQueueName() );
- }
- if (from.getDomain() != null) {
- to.setDomain( from.getDomain() );
+ /**
+ * Gets the last instance of each failed task and reschedule each Gets all cancelled tasks and
+ * schedule all of them except JOIN (join should change status to INPROGRESS) Switch workflow
+ * back to RUNNING status and call decider.
+ *
+ * @param workflowId the id of the workflow to be retried
+ */
+ public void retry(String workflowId, boolean resumeSubworkflowTasks) {
+ WorkflowModel workflow = executionDAOFacade.getWorkflowModel(workflowId, true);
+ if (!workflow.getStatus().isTerminal()) {
+ throw new ApplicationException(
+ CONFLICT, "Workflow is still running. status=" + workflow.getStatus());
+ }
+ if (workflow.getTasks().isEmpty()) {
+ throw new ApplicationException(CONFLICT, "Workflow has not started yet");
+ }
+
+ if (resumeSubworkflowTasks) {
+ Optional taskToRetry =
+ workflow.getTasks().stream().filter(UNSUCCESSFUL_TERMINAL_TASK).findFirst();
+ if (taskToRetry.isPresent()) {
+ workflow = findLastFailedSubWorkflowIfAny(taskToRetry.get(), workflow);
+ retry(workflow);
+ updateAndPushParents(workflow, "retried");
+ }
+ } else {
+ retry(workflow);
+ updateAndPushParents(workflow, "retried");
}
- if (from.getWorkerId() != null) {
- to.setWorkerId( from.getWorkerId() );
- }
- to.setLastPollTime( from.getLastPollTime() );
- return to.build();
- }
-
- public PollData fromProto(PollDataPb.PollData from) {
- PollData to = new PollData();
- to.setQueueName( from.getQueueName() );
- to.setDomain( from.getDomain() );
- to.setWorkerId( from.getWorkerId() );
- to.setLastPollTime( from.getLastPollTime() );
- return to;
}
- public RerunWorkflowRequestPb.RerunWorkflowRequest toProto(RerunWorkflowRequest from) {
- RerunWorkflowRequestPb.RerunWorkflowRequest.Builder to = RerunWorkflowRequestPb.RerunWorkflowRequest.newBuilder();
- if (from.getReRunFromWorkflowId() != null) {
- to.setReRunFromWorkflowId( from.getReRunFromWorkflowId() );
- }
- for (Map.Entry pair : from.getWorkflowInput().entrySet()) {
- to.putWorkflowInput( pair.getKey(), toProto( pair.getValue() ) );
- }
- if (from.getReRunFromTaskId() != null) {
- to.setReRunFromTaskId( from.getReRunFromTaskId() );
- }
- for (Map.Entry pair : from.getTaskInput().entrySet()) {
- to.putTaskInput( pair.getKey(), toProto( pair.getValue() ) );
+ private void updateAndPushParents(WorkflowModel workflow, String operation) {
+ String workflowIdentifier = "";
+ while (workflow.hasParent()) {
+ // update parent's sub workflow task
+ TaskModel subWorkflowTask =
+ executionDAOFacade.getTaskModel(workflow.getParentWorkflowTaskId());
+ subWorkflowTask.setSubworkflowChanged(true);
+ subWorkflowTask.setStatus(IN_PROGRESS);
+ executionDAOFacade.updateTask(subWorkflowTask);
+
+ // add an execution log
+ String currentWorkflowIdentifier = workflow.toShortString();
+ workflowIdentifier =
+ !workflowIdentifier.equals("")
+ ? String.format(
+ "%s -> %s", currentWorkflowIdentifier, workflowIdentifier)
+ : currentWorkflowIdentifier;
+ TaskExecLog log =
+ new TaskExecLog(
+ String.format("Sub workflow %s %s.", workflowIdentifier, operation));
+ log.setTaskId(subWorkflowTask.getTaskId());
+ executionDAOFacade.addTaskExecLog(Collections.singletonList(log));
+ LOGGER.info("Task {} updated. {}", log.getTaskId(), log.getLog());
+
+ // push the parent workflow to decider queue for asynchronous 'decide'
+ String parentWorkflowId = workflow.getParentWorkflowId();
+ WorkflowModel parentWorkflow =
+ executionDAOFacade.getWorkflowModel(parentWorkflowId, true);
+ parentWorkflow.setStatus(WorkflowModel.Status.RUNNING);
+ parentWorkflow.setLastRetriedTime(System.currentTimeMillis());
+ executionDAOFacade.updateWorkflow(parentWorkflow);
+ pushParentWorkflow(parentWorkflowId);
+
+ workflow = parentWorkflow;
}
- if (from.getCorrelationId() != null) {
- to.setCorrelationId( from.getCorrelationId() );
- }
- return to.build();
}
- public RerunWorkflowRequest fromProto(RerunWorkflowRequestPb.RerunWorkflowRequest from) {
- RerunWorkflowRequest to = new RerunWorkflowRequest();
- to.setReRunFromWorkflowId( from.getReRunFromWorkflowId() );
- Map workflowInputMap = new HashMap();
- for (Map.Entry pair : from.getWorkflowInputMap().entrySet()) {
- workflowInputMap.put( pair.getKey(), fromProto( pair.getValue() ) );
- }
- to.setWorkflowInput(workflowInputMap);
- to.setReRunFromTaskId( from.getReRunFromTaskId() );
- Map taskInputMap = new HashMap();
- for (Map.Entry pair : from.getTaskInputMap().entrySet()) {
- taskInputMap.put( pair.getKey(), fromProto( pair.getValue() ) );
- }
- to.setTaskInput(taskInputMap);
- to.setCorrelationId( from.getCorrelationId() );
- return to;
+ private void retry(WorkflowModel workflow) {
+ // Get all FAILED or CANCELED tasks that are not COMPLETED (or reach other terminal states)
+ // on further executions.
+ // // Eg: for Seq of tasks task1.CANCELED, task1.COMPLETED, task1 shouldn't be retried.
+ // Throw an exception if there are no FAILED tasks.
+ // Handle JOIN task CANCELED status as special case.
+ Map retriableMap = new HashMap<>();
+ for (TaskModel task : workflow.getTasks()) {
+ switch (task.getStatus()) {
+ case FAILED:
+ case FAILED_WITH_TERMINAL_ERROR:
+ case TIMED_OUT:
+ retriableMap.put(task.getReferenceTaskName(), task);
+ break;
+ case CANCELED:
+ if (task.getTaskType().equalsIgnoreCase(TaskType.JOIN.toString())
+ || task.getTaskType().equalsIgnoreCase(TaskType.DO_WHILE.toString())) {
+ task.setStatus(IN_PROGRESS);
+ // Task doesn't have to be updated yet. Will be updated along with other
+ // Workflow tasks downstream.
+ } else {
+ retriableMap.put(task.getReferenceTaskName(), task);
+ }
+ break;
+ default:
+ retriableMap.remove(task.getReferenceTaskName());
+ break;
+ }
+ }
+
+ // if workflow TIMED_OUT due to timeoutSeconds configured in the workflow definition,
+ // it may not have any unsuccessful tasks that can be retried
+ if (retriableMap.values().size() == 0
+ && workflow.getStatus() != WorkflowModel.Status.TIMED_OUT) {
+ throw new ApplicationException(
+ CONFLICT,
+ "There are no retryable tasks! Use restart if you want to attempt entire workflow execution again.");
+ }
+
+ // Update Workflow with new status.
+ // This should load Workflow from archive, if archived.
+ workflow.setStatus(WorkflowModel.Status.RUNNING);
+ workflow.setLastRetriedTime(System.currentTimeMillis());
+ String lastReasonForIncompletion = workflow.getReasonForIncompletion();
+ workflow.setReasonForIncompletion(null);
+ // Add to decider queue
+ queueDAO.push(
+ DECIDER_QUEUE,
+ workflow.getWorkflowId(),
+ workflow.getPriority(),
+ properties.getWorkflowOffsetTimeout().getSeconds());
+ executionDAOFacade.updateWorkflow(workflow);
+ LOGGER.info(
+ "Workflow {} that failed due to '{}' was retried",
+ workflow.toShortString(),
+ lastReasonForIncompletion);
+
+ // taskToBeRescheduled would set task `retried` to true, and hence it's important to
+ // updateTasks after obtaining task copy from taskToBeRescheduled.
+ final WorkflowModel finalWorkflow = workflow;
+ List retriableTasks =
+ retriableMap.values().stream()
+ .sorted(Comparator.comparingInt(TaskModel::getSeq))
+ .map(task -> taskToBeRescheduled(finalWorkflow, task))
+ .collect(Collectors.toList());
+
+ dedupAndAddTasks(workflow, retriableTasks);
+ // Note: updateTasks before updateWorkflow might fail when Workflow is archived and doesn't
+ // exist in primary store.
+ executionDAOFacade.updateTasks(workflow.getTasks());
+ scheduleTask(workflow, retriableTasks);
}
- public SkipTaskRequest fromProto(SkipTaskRequestPb.SkipTaskRequest from) {
- SkipTaskRequest to = new SkipTaskRequest();
- Map taskInputMap = new HashMap();
- for (Map.Entry pair : from.getTaskInputMap().entrySet()) {
- taskInputMap.put( pair.getKey(), fromProto( pair.getValue() ) );
- }
- to.setTaskInput(taskInputMap);
- Map taskOutputMap = new HashMap();
- for (Map.Entry pair : from.getTaskOutputMap().entrySet()) {
- taskOutputMap.put( pair.getKey(), fromProto( pair.getValue() ) );
- }
- to.setTaskOutput(taskOutputMap);
- if (from.hasTaskInputMessage()) {
- to.setTaskInputMessage( fromProto( from.getTaskInputMessage() ) );
- }
- if (from.hasTaskOutputMessage()) {
- to.setTaskOutputMessage( fromProto( from.getTaskOutputMessage() ) );
- }
- return to;
+ private WorkflowModel findLastFailedSubWorkflowIfAny(
+ TaskModel task, WorkflowModel parentWorkflow) {
+ if (TaskType.TASK_TYPE_SUB_WORKFLOW.equals(task.getTaskType())
+ && UNSUCCESSFUL_TERMINAL_TASK.test(task)) {
+ WorkflowModel subWorkflow =
+ executionDAOFacade.getWorkflowModel(task.getSubWorkflowId(), true);
+ Optional taskToRetry =
+ subWorkflow.getTasks().stream().filter(UNSUCCESSFUL_TERMINAL_TASK).findFirst();
+ if (taskToRetry.isPresent()) {
+ return findLastFailedSubWorkflowIfAny(taskToRetry.get(), subWorkflow);
+ }
+ }
+ return parentWorkflow;
}
- public StartWorkflowRequestPb.StartWorkflowRequest toProto(StartWorkflowRequest from) {
- StartWorkflowRequestPb.StartWorkflowRequest.Builder to = StartWorkflowRequestPb.StartWorkflowRequest.newBuilder();
- if (from.getName() != null) {
- to.setName( from.getName() );
- }
- if (from.getVersion() != null) {
- to.setVersion( from.getVersion() );
- }
- if (from.getCorrelationId() != null) {
- to.setCorrelationId( from.getCorrelationId() );
- }
- for (Map.Entry pair : from.getInput().entrySet()) {
- to.putInput( pair.getKey(), toProto( pair.getValue() ) );
- }
- to.putAllTaskToDomain( from.getTaskToDomain() );
- if (from.getWorkflowDef() != null) {
- to.setWorkflowDef( toProto( from.getWorkflowDef() ) );
- }
- if (from.getExternalInputPayloadStoragePath() != null) {
- to.setExternalInputPayloadStoragePath( from.getExternalInputPayloadStoragePath() );
- }
- if (from.getPriority() != null) {
- to.setPriority( from.getPriority() );
- }
- return to.build();
+ /**
+ * Reschedule a task
+ *
+ * @param task failed or cancelled task
+ * @return new instance of a task with "SCHEDULED" status
+ */
+ private TaskModel taskToBeRescheduled(WorkflowModel workflow, TaskModel task) {
+ TaskModel taskToBeRetried = task.copy();
+ taskToBeRetried.setTaskId(IDGenerator.generate());
+ taskToBeRetried.setRetriedTaskId(task.getTaskId());
+ taskToBeRetried.setStatus(SCHEDULED);
+ taskToBeRetried.setRetryCount(task.getRetryCount() + 1);
+ taskToBeRetried.setRetried(false);
+ taskToBeRetried.setPollCount(0);
+ taskToBeRetried.setCallbackAfterSeconds(0);
+ taskToBeRetried.setSubWorkflowId(null);
+ taskToBeRetried.setScheduledTime(0);
+ taskToBeRetried.setStartTime(0);
+ taskToBeRetried.setEndTime(0);
+ taskToBeRetried.setWorkerId(null);
+ taskToBeRetried.setReasonForIncompletion(null);
+ taskToBeRetried.setSeq(0);
+
+ // perform parameter replacement for retried task
+ Map taskInput =
+ parametersUtils.getTaskInput(
+ taskToBeRetried.getWorkflowTask().getInputParameters(),
+ workflow,
+ taskToBeRetried.getWorkflowTask().getTaskDefinition(),
+ taskToBeRetried.getTaskId());
+ taskToBeRetried.getInputData().putAll(taskInput);
+
+ task.setRetried(true);
+ // since this task is being retried and a retry has been computed, task lifecycle is
+ // complete
+ task.setExecuted(true);
+ return taskToBeRetried;
}
- public StartWorkflowRequest fromProto(StartWorkflowRequestPb.StartWorkflowRequest from) {
- StartWorkflowRequest to = new StartWorkflowRequest();
- to.setName( from.getName() );
- to.setVersion( from.getVersion() );
- to.setCorrelationId( from.getCorrelationId() );
- Map inputMap = new HashMap();
- for (Map.Entry pair : from.getInputMap().entrySet()) {
- inputMap.put( pair.getKey(), fromProto( pair.getValue() ) );
- }
- to.setInput(inputMap);
- to.setTaskToDomain( from.getTaskToDomainMap() );
- if (from.hasWorkflowDef()) {
- to.setWorkflowDef( fromProto( from.getWorkflowDef() ) );
- }
- to.setExternalInputPayloadStoragePath( from.getExternalInputPayloadStoragePath() );
- to.setPriority( from.getPriority() );
- return to;
+ private void endExecution(WorkflowModel workflow) {
+ Optional terminateTask =
+ workflow.getTasks().stream()
+ .filter(
+ t ->
+ TaskType.TERMINATE.name().equals(t.getTaskType())
+ && t.getStatus().isTerminal()
+ && t.getStatus().isSuccessful())
+ .findFirst();
+ if (terminateTask.isPresent()) {
+ String terminationStatus =
+ (String)
+ terminateTask
+ .get()
+ .getWorkflowTask()
+ .getInputParameters()
+ .get(Terminate.getTerminationStatusParameter());
+ String reason =
+ (String)
+ terminateTask
+ .get()
+ .getWorkflowTask()
+ .getInputParameters()
+ .get(Terminate.getTerminationReasonParameter());
+ if (StringUtils.isBlank(reason)) {
+ reason =
+ String.format(
+ "Workflow is %s by TERMINATE task: %s",
+ terminationStatus, terminateTask.get().getTaskId());
+ }
+ if (WorkflowModel.Status.FAILED.name().equals(terminationStatus)) {
+ workflow.setStatus(WorkflowModel.Status.FAILED);
+ workflow =
+ terminate(
+ workflow,
+ new TerminateWorkflowException(
+ reason, workflow.getStatus(), terminateTask.get()));
+ } else {
+ workflow.setReasonForIncompletion(reason);
+ workflow = completeWorkflow(workflow);
+ }
+ } else {
+ workflow = completeWorkflow(workflow);
+ }
+ cancelNonTerminalTasks(workflow);
}
- public SubWorkflowParamsPb.SubWorkflowParams toProto(SubWorkflowParams from) {
- SubWorkflowParamsPb.SubWorkflowParams.Builder to = SubWorkflowParamsPb.SubWorkflowParams.newBuilder();
- if (from.getName() != null) {
- to.setName( from.getName() );
- }
- if (from.getVersion() != null) {
- to.setVersion( from.getVersion() );
- }
- to.putAllTaskToDomain( from.getTaskToDomain() );
- if (from.getWorkflowDefinition() != null) {
- to.setWorkflowDefinition( toProto( from.getWorkflowDefinition() ) );
- }
- return to.build();
+ /**
+ * @param workflow the workflow to be completed
+ * @throws ApplicationException if workflow is not in terminal state
+ */
+ @VisibleForTesting
+ WorkflowModel completeWorkflow(WorkflowModel workflow) {
+ LOGGER.debug("Completing workflow execution for {}", workflow.getWorkflowId());
+
+ if (workflow.getStatus().equals(WorkflowModel.Status.COMPLETED)) {
+ queueDAO.remove(DECIDER_QUEUE, workflow.getWorkflowId()); // remove from the sweep queue
+ executionDAOFacade.removeFromPendingWorkflow(
+ workflow.getWorkflowName(), workflow.getWorkflowId());
+ LOGGER.debug("Workflow: {} has already been completed.", workflow.getWorkflowId());
+ return workflow;
+ }
+
+ if (workflow.getStatus().isTerminal()) {
+ String msg =
+ "Workflow is already in terminal state. Current status: "
+ + workflow.getStatus();
+ throw new ApplicationException(CONFLICT, msg);
+ }
+
+ // FIXME Backwards compatibility for legacy workflows already running.
+ // This code will be removed in a future version.
+ if (workflow.getWorkflowDefinition() == null) {
+ workflow = metadataMapperService.populateWorkflowWithDefinitions(workflow);
+ }
+ deciderService.updateWorkflowOutput(workflow, null);
+
+ workflow.setStatus(WorkflowModel.Status.COMPLETED);
+
+ // update the failed reference task names
+ workflow.getFailedReferenceTaskNames()
+ .addAll(
+ workflow.getTasks().stream()
+ .filter(
+ t ->
+ FAILED.equals(t.getStatus())
+ || FAILED_WITH_TERMINAL_ERROR.equals(
+ t.getStatus()))
+ .map(TaskModel::getReferenceTaskName)
+ .collect(Collectors.toSet()));
+
+ executionDAOFacade.updateWorkflow(workflow);
+ LOGGER.debug("Completed workflow execution for {}", workflow.getWorkflowId());
+ workflowStatusListener.onWorkflowCompletedIfEnabled(workflow);
+ Monitors.recordWorkflowCompletion(
+ workflow.getWorkflowName(),
+ workflow.getEndTime() - workflow.getCreateTime(),
+ workflow.getOwnerApp());
+
+ if (workflow.hasParent()) {
+ updateParentWorkflowTask(workflow);
+ LOGGER.info(
+ "{} updated parent {} task {}",
+ workflow.toShortString(),
+ workflow.getParentWorkflowId(),
+ workflow.getParentWorkflowTaskId());
+ pushParentWorkflow(workflow.getParentWorkflowId());
+ }
+
+ executionLockService.releaseLock(workflow.getWorkflowId());
+ executionLockService.deleteLock(workflow.getWorkflowId());
+ return workflow;
}
- public SubWorkflowParams fromProto(SubWorkflowParamsPb.SubWorkflowParams from) {
- SubWorkflowParams to = new SubWorkflowParams();
- to.setName( from.getName() );
- to.setVersion( from.getVersion() );
- to.setTaskToDomain( from.getTaskToDomainMap() );
- if (from.hasWorkflowDefinition()) {
- to.setWorkflowDefinition( fromProto( from.getWorkflowDefinition() ) );
+ public void terminateWorkflow(String workflowId, String reason) {
+ WorkflowModel workflow = executionDAOFacade.getWorkflowModel(workflowId, true);
+ if (WorkflowModel.Status.COMPLETED.equals(workflow.getStatus())) {
+ throw new ApplicationException(CONFLICT, "Cannot terminate a COMPLETED workflow.");
}
- return to;
+ workflow.setStatus(WorkflowModel.Status.TERMINATED);
+ terminateWorkflow(workflow, reason, null);
}
- public TaskPb.Task toProto(Task from) {
- TaskPb.Task.Builder to = TaskPb.Task.newBuilder();
- if (from.getTaskType() != null) {
- to.setTaskType( from.getTaskType() );
- }
- if (from.getStatus() != null) {
- to.setStatus( toProto( from.getStatus() ) );
- }
- for (Map.Entry pair : from.getInputData().entrySet()) {
- to.putInputData( pair.getKey(), toProto( pair.getValue() ) );
- }
- if (from.getReferenceTaskName() != null) {
- to.setReferenceTaskName( from.getReferenceTaskName() );
- }
- to.setRetryCount( from.getRetryCount() );
- to.setSeq( from.getSeq() );
- if (from.getCorrelationId() != null) {
- to.setCorrelationId( from.getCorrelationId() );
- }
- to.setPollCount( from.getPollCount() );
- if (from.getTaskDefName() != null) {
- to.setTaskDefName( from.getTaskDefName() );
- }
- to.setScheduledTime( from.getScheduledTime() );
- to.setStartTime( from.getStartTime() );
- to.setEndTime( from.getEndTime() );
- to.setUpdateTime( from.getUpdateTime() );
- to.setStartDelayInSeconds( from.getStartDelayInSeconds() );
- if (from.getRetriedTaskId() != null) {
- to.setRetriedTaskId( from.getRetriedTaskId() );
- }
- to.setRetried( from.isRetried() );
- to.setExecuted( from.isExecuted() );
- to.setCallbackFromWorker( from.isCallbackFromWorker() );
- to.setResponseTimeoutSeconds( from.getResponseTimeoutSeconds() );
- if (from.getWorkflowInstanceId() != null) {
- to.setWorkflowInstanceId( from.getWorkflowInstanceId() );
- }
- if (from.getWorkflowType() != null) {
- to.setWorkflowType( from.getWorkflowType() );
- }
- if (from.getTaskId() != null) {
- to.setTaskId( from.getTaskId() );
- }
- if (from.getReasonForIncompletion() != null) {
- to.setReasonForIncompletion( from.getReasonForIncompletion() );
+ /**
+ * @param workflow the workflow to be terminated
+ * @param reason the reason for termination
+ * @param failureWorkflow the failure workflow (if any) to be triggered as a result of this
+ * termination
+ */
+ public WorkflowModel terminateWorkflow(
+ WorkflowModel workflow, String reason, String failureWorkflow) {
+ try {
+ executionLockService.acquireLock(workflow.getWorkflowId(), 60000);
+
+ if (!workflow.getStatus().isTerminal()) {
+ workflow.setStatus(WorkflowModel.Status.TERMINATED);
+ }
+
+ // FIXME Backwards compatibility for legacy workflows already running.
+ // This code will be removed in a future version.
+ if (workflow.getWorkflowDefinition() == null) {
+ workflow = metadataMapperService.populateWorkflowWithDefinitions(workflow);
+ }
+
+ try {
+ deciderService.updateWorkflowOutput(workflow, null);
+ } catch (Exception e) {
+ // catch any failure in this step and continue the execution of terminating workflow
+ LOGGER.error(
+ "Failed to update output data for workflow: {}",
+ workflow.getWorkflowId(),
+ e);
+ Monitors.error(CLASS_NAME, "terminateWorkflow");
+ }
+
+ // update the failed reference task names
+ workflow.getFailedReferenceTaskNames()
+ .addAll(
+ workflow.getTasks().stream()
+ .filter(
+ t ->
+ FAILED.equals(t.getStatus())
+ || FAILED_WITH_TERMINAL_ERROR.equals(
+ t.getStatus()))
+ .map(TaskModel::getReferenceTaskName)
+ .collect(Collectors.toSet()));
+
+ String workflowId = workflow.getWorkflowId();
+ workflow.setReasonForIncompletion(reason);
+ executionDAOFacade.updateWorkflow(workflow);
+ workflowStatusListener.onWorkflowTerminatedIfEnabled(workflow);
+ Monitors.recordWorkflowTermination(
+ workflow.getWorkflowName(), workflow.getStatus(), workflow.getOwnerApp());
+
+ List tasks = workflow.getTasks();
+ try {
+ // Remove from the task queue if they were there
+ tasks.forEach(
+ task -> queueDAO.remove(QueueUtils.getQueueName(task), task.getTaskId()));
+ } catch (Exception e) {
+ LOGGER.warn(
+ "Error removing task(s) from queue during workflow termination : {}",
+ workflowId,
+ e);
+ }
+
+ if (workflow.hasParent()) {
+ updateParentWorkflowTask(workflow);
+ LOGGER.info(
+ "{} updated parent {} task {}",
+ workflow.toShortString(),
+ workflow.getParentWorkflowId(),
+ workflow.getParentWorkflowTaskId());
+ pushParentWorkflow(workflow.getParentWorkflowId());
+ }
+
+ if (!StringUtils.isBlank(failureWorkflow)) {
+ Map input = new HashMap<>(workflow.getInput());
+ input.put("workflowId", workflowId);
+ input.put("reason", reason);
+ input.put("failureStatus", workflow.getStatus().toString());
+ if (workflow.getFailedTaskId() != null) {
+ input.put("failureTaskId", workflow.getFailedTaskId());
+ }
+
+ try {
+ WorkflowDef latestFailureWorkflow =
+ metadataDAO
+ .getLatestWorkflowDef(failureWorkflow)
+ .orElseThrow(
+ () ->
+ new RuntimeException(
+ "Failure Workflow Definition not found for: "
+ + failureWorkflow));
+
+ String failureWFId =
+ startWorkflow(
+ latestFailureWorkflow,
+ input,
+ null,
+ workflow.getCorrelationId(),
+ null,
+ workflow.getTaskToDomain());
+
+ workflow.getOutput().put("conductor.failure_workflow", failureWFId);
+ } catch (Exception e) {
+ LOGGER.error("Failed to start error workflow", e);
+ workflow.getOutput()
+ .put(
+ "conductor.failure_workflow",
+ "Error workflow "
+ + failureWorkflow
+ + " failed to start. reason: "
+ + e.getMessage());
+ Monitors.recordWorkflowStartError(
+ failureWorkflow, WorkflowContext.get().getClientApp());
+ }
+ executionDAOFacade.updateWorkflow(workflow);
+ }
+ executionDAOFacade.removeFromPendingWorkflow(
+ workflow.getWorkflowName(), workflow.getWorkflowId());
+
+ List erroredTasks = cancelNonTerminalTasks(workflow);
+ if (!erroredTasks.isEmpty()) {
+ throw new ApplicationException(
+ Code.INTERNAL_ERROR,
+ String.format(
+ "Error canceling system tasks: %s",
+ String.join(",", erroredTasks)));
+ }
+ return workflow;
+ } finally {
+ executionLockService.releaseLock(workflow.getWorkflowId());
+ executionLockService.deleteLock(workflow.getWorkflowId());
}
- to.setCallbackAfterSeconds( from.getCallbackAfterSeconds() );
- if (from.getWorkerId() != null) {
- to.setWorkerId( from.getWorkerId() );
- }
- for (Map.Entry pair : from.getOutputData().entrySet()) {
- to.putOutputData( pair.getKey(), toProto( pair.getValue() ) );
- }
- if (from.getWorkflowTask() != null) {
- to.setWorkflowTask( toProto( from.getWorkflowTask() ) );
- }
- if (from.getDomain() != null) {
- to.setDomain( from.getDomain() );
- }
- if (from.getInputMessage() != null) {
- to.setInputMessage( toProto( from.getInputMessage() ) );
- }
- if (from.getOutputMessage() != null) {
- to.setOutputMessage( toProto( from.getOutputMessage() ) );
- }
- to.setRateLimitPerFrequency( from.getRateLimitPerFrequency() );
- to.setRateLimitFrequencyInSeconds( from.getRateLimitFrequencyInSeconds() );
- if (from.getExternalInputPayloadStoragePath() != null) {
- to.setExternalInputPayloadStoragePath( from.getExternalInputPayloadStoragePath() );
- }
- if (from.getExternalOutputPayloadStoragePath() != null) {
- to.setExternalOutputPayloadStoragePath( from.getExternalOutputPayloadStoragePath() );
- }
- to.setWorkflowPriority( from.getWorkflowPriority() );
- if (from.getExecutionNameSpace() != null) {
- to.setExecutionNameSpace( from.getExecutionNameSpace() );
- }
- if (from.getIsolationGroupId() != null) {
- to.setIsolationGroupId( from.getIsolationGroupId() );
- }
- to.setIteration( from.getIteration() );
- if (from.getSubWorkflowId() != null) {
- to.setSubWorkflowId( from.getSubWorkflowId() );
- }
- to.setSubworkflowChanged( from.isSubworkflowChanged() );
- return to.build();
}
- public Task fromProto(TaskPb.Task from) {
- Task to = new Task();
- to.setTaskType( from.getTaskType() );
- to.setStatus( fromProto( from.getStatus() ) );
- Map inputDataMap = new HashMap();
- for (Map.Entry pair : from.getInputDataMap().entrySet()) {
- inputDataMap.put( pair.getKey(), fromProto( pair.getValue() ) );
- }
- to.setInputData(inputDataMap);
- to.setReferenceTaskName( from.getReferenceTaskName() );
- to.setRetryCount( from.getRetryCount() );
- to.setSeq( from.getSeq() );
- to.setCorrelationId( from.getCorrelationId() );
- to.setPollCount( from.getPollCount() );
- to.setTaskDefName( from.getTaskDefName() );
- to.setScheduledTime( from.getScheduledTime() );
- to.setStartTime( from.getStartTime() );
- to.setEndTime( from.getEndTime() );
- to.setUpdateTime( from.getUpdateTime() );
- to.setStartDelayInSeconds( from.getStartDelayInSeconds() );
- to.setRetriedTaskId( from.getRetriedTaskId() );
- to.setRetried( from.getRetried() );
- to.setExecuted( from.getExecuted() );
- to.setCallbackFromWorker( from.getCallbackFromWorker() );
- to.setResponseTimeoutSeconds( from.getResponseTimeoutSeconds() );
- to.setWorkflowInstanceId( from.getWorkflowInstanceId() );
- to.setWorkflowType( from.getWorkflowType() );
- to.setTaskId( from.getTaskId() );
- to.setReasonForIncompletion( from.getReasonForIncompletion() );
- to.setCallbackAfterSeconds( from.getCallbackAfterSeconds() );
- to.setWorkerId( from.getWorkerId() );
- Map outputDataMap = new HashMap();
- for (Map.Entry pair : from.getOutputDataMap().entrySet()) {
- outputDataMap.put( pair.getKey(), fromProto( pair.getValue() ) );
- }
- to.setOutputData(outputDataMap);
- if (from.hasWorkflowTask()) {
- to.setWorkflowTask( fromProto( from.getWorkflowTask() ) );
- }
- to.setDomain( from.getDomain() );
- if (from.hasInputMessage()) {
- to.setInputMessage( fromProto( from.getInputMessage() ) );
- }
- if (from.hasOutputMessage()) {
- to.setOutputMessage( fromProto( from.getOutputMessage() ) );
- }
- to.setRateLimitPerFrequency( from.getRateLimitPerFrequency() );
- to.setRateLimitFrequencyInSeconds( from.getRateLimitFrequencyInSeconds() );
- to.setExternalInputPayloadStoragePath( from.getExternalInputPayloadStoragePath() );
- to.setExternalOutputPayloadStoragePath( from.getExternalOutputPayloadStoragePath() );
- to.setWorkflowPriority( from.getWorkflowPriority() );
- to.setExecutionNameSpace( from.getExecutionNameSpace() );
- to.setIsolationGroupId( from.getIsolationGroupId() );
- to.setIteration( from.getIteration() );
- to.setSubWorkflowId( from.getSubWorkflowId() );
- to.setSubworkflowChanged( from.getSubworkflowChanged() );
- return to;
+ /**
+ * @param taskResult the task result to be updated
+ * @throws ApplicationException
+ */
+ public void updateTask(TaskResult taskResult) {
+ if (taskResult == null) {
+ throw new ApplicationException(
+ ApplicationException.Code.INVALID_INPUT, "Task object is null");
+ }
+
+ String workflowId = taskResult.getWorkflowInstanceId();
+ WorkflowModel workflowInstance = executionDAOFacade.getWorkflowModel(workflowId, true);
+
+ // FIXME Backwards compatibility for legacy workflows already running.
+ // This code will be removed in a future version.
+ if (workflowInstance.getWorkflowDefinition() == null) {
+ workflowInstance =
+ metadataMapperService.populateWorkflowWithDefinitions(workflowInstance);
+ }
+
+ TaskModel task =
+ Optional.ofNullable(executionDAOFacade.getTaskModel(taskResult.getTaskId()))
+ .orElseThrow(
+ () ->
+ new ApplicationException(
+ ApplicationException.Code.NOT_FOUND,
+ "No such task found by id: "
+ + taskResult.getTaskId()));
+
+ LOGGER.debug("Task: {} belonging to Workflow {} being updated", task, workflowInstance);
+
+ String taskQueueName = QueueUtils.getQueueName(task);
+
+ if (task.getStatus().isTerminal()) {
+ // Task was already updated....
+ queueDAO.remove(taskQueueName, taskResult.getTaskId());
+ LOGGER.info(
+ "Task: {} has already finished execution with status: {} within workflow: {}. Removed task from queue: {}",
+ task.getTaskId(),
+ task.getStatus(),
+ task.getWorkflowInstanceId(),
+ taskQueueName);
+ Monitors.recordUpdateConflict(
+ task.getTaskType(), workflowInstance.getWorkflowName(), task.getStatus());
+ return;
+ }
+
+ if (workflowInstance.getStatus().isTerminal()) {
+ // Workflow is in terminal state
+ queueDAO.remove(taskQueueName, taskResult.getTaskId());
+ LOGGER.info(
+ "Workflow: {} has already finished execution. Task update for: {} ignored and removed from Queue: {}.",
+ workflowInstance,
+ taskResult.getTaskId(),
+ taskQueueName);
+ Monitors.recordUpdateConflict(
+ task.getTaskType(),
+ workflowInstance.getWorkflowName(),
+ workflowInstance.getStatus());
+ return;
+ }
+
+ // for system tasks, setting to SCHEDULED would mean restarting the task which is
+ // undesirable
+ // for worker tasks, set status to SCHEDULED and push to the queue
+ if (!systemTaskRegistry.isSystemTask(task.getTaskType())
+ && taskResult.getStatus() == TaskResult.Status.IN_PROGRESS) {
+ task.setStatus(SCHEDULED);
+ } else {
+ task.setStatus(TaskModel.Status.valueOf(taskResult.getStatus().name()));
+ }
+ task.setOutputMessage(taskResult.getOutputMessage());
+ task.setReasonForIncompletion(taskResult.getReasonForIncompletion());
+ task.setWorkerId(taskResult.getWorkerId());
+ task.setCallbackAfterSeconds(taskResult.getCallbackAfterSeconds());
+ task.setOutputData(taskResult.getOutputData());
+ task.setSubWorkflowId(taskResult.getSubWorkflowId());
+
+ if (StringUtils.isNotBlank(taskResult.getExternalOutputPayloadStoragePath())) {
+ task.setExternalOutputPayloadStoragePath(
+ taskResult.getExternalOutputPayloadStoragePath());
+ }
+
+ if (task.getStatus().isTerminal()) {
+ task.setEndTime(System.currentTimeMillis());
+ }
+
+ // Update message in Task queue based on Task status
+ switch (task.getStatus()) {
+ case COMPLETED:
+ case CANCELED:
+ case FAILED:
+ case FAILED_WITH_TERMINAL_ERROR:
+ case TIMED_OUT:
+ try {
+ queueDAO.remove(taskQueueName, taskResult.getTaskId());
+ LOGGER.debug(
+ "Task: {} removed from taskQueue: {} since the task status is {}",
+ task,
+ taskQueueName,
+ task.getStatus().name());
+ } catch (Exception e) {
+ // Ignore exceptions on queue remove as it wouldn't impact task and workflow
+ // execution, and will be cleaned up eventually
+ String errorMsg =
+ String.format(
+ "Error removing the message in queue for task: %s for workflow: %s",
+ task.getTaskId(), workflowId);
+ LOGGER.warn(errorMsg, e);
+ Monitors.recordTaskQueueOpError(
+ task.getTaskType(), workflowInstance.getWorkflowName());
+ }
+ break;
+ case IN_PROGRESS:
+ case SCHEDULED:
+ try {
+ String postponeTaskMessageDesc =
+ "Postponing Task message in queue for taskId: " + task.getTaskId();
+ String postponeTaskMessageOperation = "postponeTaskMessage";
+
+ new RetryUtil<>()
+ .retryOnException(
+ () -> {
+ // postpone based on callbackAfterSeconds
+ long callBack = taskResult.getCallbackAfterSeconds();
+ queueDAO.postpone(
+ taskQueueName,
+ task.getTaskId(),
+ task.getWorkflowPriority(),
+ callBack);
+ LOGGER.debug(
+ "Task: {} postponed in taskQueue: {} since the task status is {} with callbackAfterSeconds: {}",
+ task,
+ taskQueueName,
+ task.getStatus().name(),
+ callBack);
+ return null;
+ },
+ null,
+ null,
+ 2,
+ postponeTaskMessageDesc,
+ postponeTaskMessageOperation);
+ } catch (Exception e) {
+ // Throw exceptions on queue postpone, this would impact task execution
+ String errorMsg =
+ String.format(
+ "Error postponing the message in queue for task: %s for workflow: %s",
+ task.getTaskId(), workflowId);
+ LOGGER.error(errorMsg, e);
+ Monitors.recordTaskQueueOpError(
+ task.getTaskType(), workflowInstance.getWorkflowName());
+ throw new ApplicationException(ApplicationException.Code.BACKEND_ERROR, e);
+ }
+ break;
+ default:
+ break;
+ }
+
+ // Throw an ApplicationException if below operations fail to avoid workflow inconsistencies.
+ try {
+ String updateTaskDesc = "Updating Task with taskId: " + task.getTaskId();
+ String updateTaskOperation = "updateTask";
+
+ new RetryUtil<>()
+ .retryOnException(
+ () -> {
+ executionDAOFacade.updateTask(task);
+ return null;
+ },
+ null,
+ null,
+ 2,
+ updateTaskDesc,
+ updateTaskOperation);
+ } catch (Exception e) {
+ String errorMsg =
+ String.format(
+ "Error updating task: %s for workflow: %s",
+ task.getTaskId(), workflowId);
+ LOGGER.error(errorMsg, e);
+ Monitors.recordTaskUpdateError(task.getTaskType(), workflowInstance.getWorkflowName());
+ throw new ApplicationException(ApplicationException.Code.BACKEND_ERROR, e);
+ }
+
+ taskResult.getLogs().forEach(taskExecLog -> taskExecLog.setTaskId(task.getTaskId()));
+ executionDAOFacade.addTaskExecLog(taskResult.getLogs());
+
+ if (task.getStatus().isTerminal()) {
+ long duration = getTaskDuration(0, task);
+ long lastDuration = task.getEndTime() - task.getStartTime();
+ Monitors.recordTaskExecutionTime(
+ task.getTaskDefName(), duration, true, task.getStatus());
+ Monitors.recordTaskExecutionTime(
+ task.getTaskDefName(), lastDuration, false, task.getStatus());
+ }
+
+ decide(workflowId);
}
- public TaskPb.Task.Status toProto(Task.Status from) {
- TaskPb.Task.Status to;
- switch (from) {
- case IN_PROGRESS: to = TaskPb.Task.Status.IN_PROGRESS; break;
- case CANCELED: to = TaskPb.Task.Status.CANCELED; break;
- case FAILED: to = TaskPb.Task.Status.FAILED; break;
- case FAILED_WITH_TERMINAL_ERROR: to = TaskPb.Task.Status.FAILED_WITH_TERMINAL_ERROR; break;
- case COMPLETED: to = TaskPb.Task.Status.COMPLETED; break;
- case COMPLETED_WITH_ERRORS: to = TaskPb.Task.Status.COMPLETED_WITH_ERRORS; break;
- case SCHEDULED: to = TaskPb.Task.Status.SCHEDULED; break;
- case TIMED_OUT: to = TaskPb.Task.Status.TIMED_OUT; break;
- case SKIPPED: to = TaskPb.Task.Status.SKIPPED; break;
- default: throw new IllegalArgumentException("Unexpected enum constant: " + from);
- }
- return to;
+ public TaskModel getTask(String taskId) {
+ return Optional.ofNullable(executionDAOFacade.getTaskModel(taskId))
+ .map(
+ task -> {
+ if (task.getWorkflowTask() != null) {
+ return metadataMapperService.populateTaskWithDefinition(task);
+ }
+ return task;
+ })
+ .orElse(null);
}
- public Task.Status fromProto(TaskPb.Task.Status from) {
- Task.Status to;
- switch (from) {
- case IN_PROGRESS: to = Task.Status.IN_PROGRESS; break;
- case CANCELED: to = Task.Status.CANCELED; break;
- case FAILED: to = Task.Status.FAILED; break;
- case FAILED_WITH_TERMINAL_ERROR: to = Task.Status.FAILED_WITH_TERMINAL_ERROR; break;
- case COMPLETED: to = Task.Status.COMPLETED; break;
- case COMPLETED_WITH_ERRORS: to = Task.Status.COMPLETED_WITH_ERRORS; break;
- case SCHEDULED: to = Task.Status.SCHEDULED; break;
- case TIMED_OUT: to = Task.Status.TIMED_OUT; break;
- case SKIPPED: to = Task.Status.SKIPPED; break;
- default: throw new IllegalArgumentException("Unexpected enum constant: " + from);
- }
- return to;
+ public List getRunningWorkflows(String workflowName, int version) {
+ return executionDAOFacade.getPendingWorkflowsByName(workflowName, version);
}
- public TaskDefPb.TaskDef toProto(TaskDef from) {
- TaskDefPb.TaskDef.Builder to = TaskDefPb.TaskDef.newBuilder();
- if (from.getName() != null) {
- to.setName( from.getName() );
- }
- if (from.getDescription() != null) {
- to.setDescription( from.getDescription() );
- }
- to.setRetryCount( from.getRetryCount() );
- to.setTimeoutSeconds( from.getTimeoutSeconds() );
- to.addAllInputKeys( from.getInputKeys() );
- to.addAllOutputKeys( from.getOutputKeys() );
- if (from.getTimeoutPolicy() != null) {
- to.setTimeoutPolicy( toProto( from.getTimeoutPolicy() ) );
- }
- if (from.getRetryLogic() != null) {
- to.setRetryLogic( toProto( from.getRetryLogic() ) );
- }
- to.setRetryDelaySeconds( from.getRetryDelaySeconds() );
- to.setResponseTimeoutSeconds( from.getResponseTimeoutSeconds() );
- if (from.getConcurrentExecLimit() != null) {
- to.setConcurrentExecLimit( from.getConcurrentExecLimit() );
- }
- for (Map.Entry pair : from.getInputTemplate().entrySet()) {
- to.putInputTemplate( pair.getKey(), toProto( pair.getValue() ) );
- }
- if (from.getRateLimitPerFrequency() != null) {
- to.setRateLimitPerFrequency( from.getRateLimitPerFrequency() );
- }
- if (from.getRateLimitFrequencyInSeconds() != null) {
- to.setRateLimitFrequencyInSeconds( from.getRateLimitFrequencyInSeconds() );
- }
- if (from.getIsolationGroupId() != null) {
- to.setIsolationGroupId( from.getIsolationGroupId() );
- }
- if (from.getExecutionNameSpace() != null) {
- to.setExecutionNameSpace( from.getExecutionNameSpace() );
- }
- if (from.getOwnerEmail() != null) {
- to.setOwnerEmail( from.getOwnerEmail() );
- }
- if (from.getPollTimeoutSeconds() != null) {
- to.setPollTimeoutSeconds( from.getPollTimeoutSeconds() );
- }
- if (from.getBackoffScaleFactor() != null) {
- to.setBackoffScaleFactor( from.getBackoffScaleFactor() );
- }
- return to.build();
+ public List getWorkflows(String name, Integer version, Long startTime, Long endTime) {
+ return executionDAOFacade.getWorkflowsByName(name, startTime, endTime).stream()
+ .filter(workflow -> workflow.getWorkflowVersion() == version)
+ .map(Workflow::getWorkflowId)
+ .collect(Collectors.toList());
}
- public TaskDef fromProto(TaskDefPb.TaskDef from) {
- TaskDef to = new TaskDef();
- to.setName( from.getName() );
- to.setDescription( from.getDescription() );
- to.setRetryCount( from.getRetryCount() );
- to.setTimeoutSeconds( from.getTimeoutSeconds() );
- to.setInputKeys( from.getInputKeysList().stream().collect(Collectors.toCollection(ArrayList::new)) );
- to.setOutputKeys( from.getOutputKeysList().stream().collect(Collectors.toCollection(ArrayList::new)) );
- to.setTimeoutPolicy( fromProto( from.getTimeoutPolicy() ) );
- to.setRetryLogic( fromProto( from.getRetryLogic() ) );
- to.setRetryDelaySeconds( from.getRetryDelaySeconds() );
- to.setResponseTimeoutSeconds( from.getResponseTimeoutSeconds() );
- to.setConcurrentExecLimit( from.getConcurrentExecLimit() );
- Map inputTemplateMap = new HashMap();
- for (Map.Entry pair : from.getInputTemplateMap().entrySet()) {
- inputTemplateMap.put( pair.getKey(), fromProto( pair.getValue() ) );
- }
- to.setInputTemplate(inputTemplateMap);
- to.setRateLimitPerFrequency( from.getRateLimitPerFrequency() );
- to.setRateLimitFrequencyInSeconds( from.getRateLimitFrequencyInSeconds() );
- to.setIsolationGroupId( from.getIsolationGroupId() );
- to.setExecutionNameSpace( from.getExecutionNameSpace() );
- to.setOwnerEmail( from.getOwnerEmail() );
- to.setPollTimeoutSeconds( from.getPollTimeoutSeconds() );
- to.setBackoffScaleFactor( from.getBackoffScaleFactor() );
- return to;
+ public List getRunningWorkflowIds(String workflowName, int version) {
+ return executionDAOFacade.getRunningWorkflowIds(workflowName, version);
}
- public TaskDefPb.TaskDef.RetryLogic toProto(TaskDef.RetryLogic from) {
- TaskDefPb.TaskDef.RetryLogic to;
- switch (from) {
- case FIXED: to = TaskDefPb.TaskDef.RetryLogic.FIXED; break;
- case EXPONENTIAL_BACKOFF: to = TaskDefPb.TaskDef.RetryLogic.EXPONENTIAL_BACKOFF; break;
- case LINEAR_BACKOFF: to = TaskDefPb.TaskDef.RetryLogic.LINEAR_BACKOFF; break;
- default: throw new IllegalArgumentException("Unexpected enum constant: " + from);
- }
- return to;
+ /**
+ * @param workflowId ID of the workflow to evaluate the state for
+ * @return true if the workflow has completed (success or failed), false otherwise.
+ * @throws ApplicationException If there was an error - caller should retry in this case.
+ */
+ public boolean decide(String workflowId) {
+ if (!executionLockService.acquireLock(workflowId)) {
+ return false;
+ }
+
+ // If it is a new workflow, the tasks will be still empty even though include tasks is true
+ WorkflowModel workflow = executionDAOFacade.getWorkflowModel(workflowId, true);
+
+ // FIXME Backwards compatibility for legacy workflows already running.
+ // This code will be removed in a future version.
+ workflow = metadataMapperService.populateWorkflowWithDefinitions(workflow);
+
+ if (workflow.getStatus().isTerminal()) {
+ if (!workflow.getStatus().isSuccessful()) {
+ cancelNonTerminalTasks(workflow);
+ }
+ return true;
+ }
+
+ // we find any sub workflow tasks that have changed
+ // and change the workflow/task state accordingly
+ adjustStateIfSubWorkflowChanged(workflow);
+
+ try {
+ DeciderService.DeciderOutcome outcome = deciderService.decide(workflow);
+ if (outcome.isComplete) {
+ endExecution(workflow);
+ return true;
+ }
+
+ List tasksToBeScheduled = outcome.tasksToBeScheduled;
+ setTaskDomains(tasksToBeScheduled, workflow);
+ List tasksToBeUpdated = outcome.tasksToBeUpdated;
+ boolean stateChanged = false;
+
+ tasksToBeScheduled = dedupAndAddTasks(workflow, tasksToBeScheduled);
+
+ for (TaskModel task : outcome.tasksToBeScheduled) {
+ if (systemTaskRegistry.isSystemTask(task.getTaskType())
+ && NON_TERMINAL_TASK.test(task)) {
+ WorkflowSystemTask workflowSystemTask =
+ systemTaskRegistry.get(task.getTaskType());
+ if (!workflowSystemTask.isAsync()
+ && workflowSystemTask.execute(workflow, task, this)) {
+ tasksToBeUpdated.add(task);
+ stateChanged = true;
+ }
+ }
+ }
+
+ if (!outcome.tasksToBeUpdated.isEmpty() || !tasksToBeScheduled.isEmpty()) {
+ executionDAOFacade.updateTasks(tasksToBeUpdated);
+ executionDAOFacade.updateWorkflow(workflow);
+ }
+
+ stateChanged = scheduleTask(workflow, tasksToBeScheduled) || stateChanged;
+
+ if (stateChanged) {
+ decide(workflowId);
+ }
+
+ } catch (TerminateWorkflowException twe) {
+ LOGGER.info("Execution terminated of workflow: {}", workflowId, twe);
+ terminate(workflow, twe);
+ return true;
+ } catch (RuntimeException e) {
+ LOGGER.error("Error deciding workflow: {}", workflowId, e);
+ throw e;
+ } finally {
+ executionLockService.releaseLock(workflowId);
+ }
+ return false;
}
- public TaskDef.RetryLogic fromProto(TaskDefPb.TaskDef.RetryLogic from) {
- TaskDef.RetryLogic to;
- switch (from) {
- case FIXED: to = TaskDef.RetryLogic.FIXED; break;
- case EXPONENTIAL_BACKOFF: to = TaskDef.RetryLogic.EXPONENTIAL_BACKOFF; break;
- case LINEAR_BACKOFF: to = TaskDef.RetryLogic.LINEAR_BACKOFF; break;
- default: throw new IllegalArgumentException("Unexpected enum constant: " + from);
+ private void adjustStateIfSubWorkflowChanged(WorkflowModel workflow) {
+ Optional changedSubWorkflowTask = findChangedSubWorkflowTask(workflow);
+ if (changedSubWorkflowTask.isPresent()) {
+ // reset the flag
+ TaskModel subWorkflowTask = changedSubWorkflowTask.get();
+ subWorkflowTask.setSubworkflowChanged(false);
+ executionDAOFacade.updateTask(subWorkflowTask);
+
+ LOGGER.info(
+ "{} reset subworkflowChanged flag for {}",
+ workflow.toShortString(),
+ subWorkflowTask.getTaskId());
+
+ // find all terminal and unsuccessful JOIN tasks and set them to IN_PROGRESS
+ if (workflow.getWorkflowDefinition().containsType(TaskType.TASK_TYPE_JOIN)
+ || workflow.getWorkflowDefinition()
+ .containsType(TaskType.TASK_TYPE_FORK_JOIN_DYNAMIC)) {
+ // if we are here, then the SUB_WORKFLOW task could be part of a FORK_JOIN or
+ // FORK_JOIN_DYNAMIC
+ // and the JOIN task(s) needs to be evaluated again, set them to IN_PROGRESS
+ workflow.getTasks().stream()
+ .filter(UNSUCCESSFUL_JOIN_TASK)
+ .peek(t -> t.setStatus(TaskModel.Status.IN_PROGRESS))
+ .forEach(executionDAOFacade::updateTask);
+ }
}
- return to;
}
- public TaskDefPb.TaskDef.TimeoutPolicy toProto(TaskDef.TimeoutPolicy from) {
- TaskDefPb.TaskDef.TimeoutPolicy to;
- switch (from) {
- case RETRY: to = TaskDefPb.TaskDef.TimeoutPolicy.RETRY; break;
- case TIME_OUT_WF: to = TaskDefPb.TaskDef.TimeoutPolicy.TIME_OUT_WF; break;
- case ALERT_ONLY: to = TaskDefPb.TaskDef.TimeoutPolicy.ALERT_ONLY; break;
- default: throw new IllegalArgumentException("Unexpected enum constant: " + from);
- }
- return to;
+ private Optional findChangedSubWorkflowTask(WorkflowModel workflow) {
+ WorkflowDef workflowDef =
+ Optional.ofNullable(workflow.getWorkflowDefinition())
+ .orElseGet(
+ () ->
+ metadataDAO
+ .getWorkflowDef(
+ workflow.getWorkflowName(),
+ workflow.getWorkflowVersion())
+ .orElseThrow(
+ () ->
+ new ApplicationException(
+ BACKEND_ERROR,
+ "Workflow Definition is not found")));
+ if (workflowDef.containsType(TaskType.TASK_TYPE_SUB_WORKFLOW)
+ || workflow.getWorkflowDefinition()
+ .containsType(TaskType.TASK_TYPE_FORK_JOIN_DYNAMIC)) {
+ return workflow.getTasks().stream()
+ .filter(
+ t ->
+ t.getTaskType().equals(TaskType.TASK_TYPE_SUB_WORKFLOW)
+ && t.isSubworkflowChanged()
+ && !t.isRetried())
+ .findFirst();
+ }
+ return Optional.empty();
}
- public TaskDef.TimeoutPolicy fromProto(TaskDefPb.TaskDef.TimeoutPolicy from) {
- TaskDef.TimeoutPolicy to;
- switch (from) {
- case RETRY: to = TaskDef.TimeoutPolicy.RETRY; break;
- case TIME_OUT_WF: to = TaskDef.TimeoutPolicy.TIME_OUT_WF; break;
- case ALERT_ONLY: to = TaskDef.TimeoutPolicy.ALERT_ONLY; break;
- default: throw new IllegalArgumentException("Unexpected enum constant: " + from);
- }
- return to;
+ @VisibleForTesting
+ List cancelNonTerminalTasks(WorkflowModel workflow) {
+ List erroredTasks = new ArrayList<>();
+ // Update non-terminal tasks' status to CANCELED
+ for (TaskModel task : workflow.getTasks()) {
+ if (!task.getStatus().isTerminal()) {
+ // Cancel the ones which are not completed yet....
+ task.setStatus(CANCELED);
+ if (systemTaskRegistry.isSystemTask(task.getTaskType())) {
+ WorkflowSystemTask workflowSystemTask =
+ systemTaskRegistry.get(task.getTaskType());
+ try {
+ workflowSystemTask.cancel(workflow, task, this);
+ } catch (Exception e) {
+ erroredTasks.add(task.getReferenceTaskName());
+ LOGGER.error(
+ "Error canceling system task:{}/{} in workflow: {}",
+ workflowSystemTask.getTaskType(),
+ task.getTaskId(),
+ workflow.getWorkflowId(),
+ e);
+ }
+ }
+ executionDAOFacade.updateTask(task);
+ }
+ }
+ if (erroredTasks.isEmpty()) {
+ try {
+ workflowStatusListener.onWorkflowFinalizedIfEnabled(workflow);
+ queueDAO.remove(DECIDER_QUEUE, workflow.getWorkflowId());
+ } catch (Exception e) {
+ LOGGER.error(
+ "Error removing workflow: {} from decider queue",
+ workflow.getWorkflowId(),
+ e);
+ }
+ }
+ return erroredTasks;
}
- public TaskExecLogPb.TaskExecLog toProto(TaskExecLog from) {
- TaskExecLogPb.TaskExecLog.Builder to = TaskExecLogPb.TaskExecLog.newBuilder();
- if (from.getLog() != null) {
- to.setLog( from.getLog() );
- }
- if (from.getTaskId() != null) {
- to.setTaskId( from.getTaskId() );
- }
- to.setCreatedTime( from.getCreatedTime() );
- return to.build();
+ @VisibleForTesting
+ List dedupAndAddTasks(WorkflowModel workflow, List tasks) {
+ List tasksInWorkflow =
+ workflow.getTasks().stream()
+ .map(task -> task.getReferenceTaskName() + "_" + task.getRetryCount())
+ .collect(Collectors.toList());
+
+ List dedupedTasks =
+ tasks.stream()
+ .filter(
+ task ->
+ !tasksInWorkflow.contains(
+ task.getReferenceTaskName()
+ + "_"
+ + task.getRetryCount()))
+ .collect(Collectors.toList());
+
+ workflow.getTasks().addAll(dedupedTasks);
+ return dedupedTasks;
}
- public TaskExecLog fromProto(TaskExecLogPb.TaskExecLog from) {
- TaskExecLog to = new TaskExecLog();
- to.setLog( from.getLog() );
- to.setTaskId( from.getTaskId() );
- to.setCreatedTime( from.getCreatedTime() );
- return to;
+ /** @throws ApplicationException if the workflow cannot be paused */
+ public void pauseWorkflow(String workflowId) {
+ try {
+ executionLockService.acquireLock(workflowId, 60000);
+ WorkflowModel.Status status = WorkflowModel.Status.PAUSED;
+ WorkflowModel workflow = executionDAOFacade.getWorkflowModel(workflowId, false);
+ if (workflow.getStatus().isTerminal()) {
+ throw new ApplicationException(
+ CONFLICT,
+ "Workflow id " + workflowId + " has ended, status cannot be updated.");
+ }
+ if (workflow.getStatus().equals(status)) {
+ return; // Already paused!
+ }
+ workflow.setStatus(status);
+ executionDAOFacade.updateWorkflow(workflow);
+ } finally {
+ executionLockService.releaseLock(workflowId);
+ }
+
+ // remove from the sweep queue
+ // any exceptions can be ignored, as this is not critical to the pause operation
+ try {
+ queueDAO.remove(DECIDER_QUEUE, workflowId);
+ } catch (Exception e) {
+ LOGGER.info(
+ "[pauseWorkflow] Error removing workflow: {} from decider queue",
+ workflowId,
+ e);
+ }
}
- public TaskResultPb.TaskResult toProto(TaskResult from) {
- TaskResultPb.TaskResult.Builder to = TaskResultPb.TaskResult.newBuilder();
- if (from.getWorkflowInstanceId() != null) {
- to.setWorkflowInstanceId( from.getWorkflowInstanceId() );
- }
- if (from.getTaskId() != null) {
- to.setTaskId( from.getTaskId() );
- }
- if (from.getReasonForIncompletion() != null) {
- to.setReasonForIncompletion( from.getReasonForIncompletion() );
- }
- to.setCallbackAfterSeconds( from.getCallbackAfterSeconds() );
- if (from.getWorkerId() != null) {
- to.setWorkerId( from.getWorkerId() );
- }
- if (from.getStatus() != null) {
- to.setStatus( toProto( from.getStatus() ) );
- }
- for (Map.Entry pair : from.getOutputData().entrySet()) {
- to.putOutputData( pair.getKey(), toProto( pair.getValue() ) );
- }
- if (from.getOutputMessage() != null) {
- to.setOutputMessage( toProto( from.getOutputMessage() ) );
- }
- return to.build();
+ /**
+ * @param workflowId the workflow to be resumed
+ * @throws IllegalStateException if the workflow is not in PAUSED state
+ */
+ public void resumeWorkflow(String workflowId) {
+ WorkflowModel workflow = executionDAOFacade.getWorkflowModel(workflowId, false);
+ if (!workflow.getStatus().equals(WorkflowModel.Status.PAUSED)) {
+ throw new IllegalStateException(
+ "The workflow "
+ + workflowId
+ + " is not PAUSED so cannot resume. "
+ + "Current status is "
+ + workflow.getStatus().name());
+ }
+ workflow.setStatus(WorkflowModel.Status.RUNNING);
+ workflow.setLastRetriedTime(System.currentTimeMillis());
+ // Add to decider queue
+ queueDAO.push(
+ DECIDER_QUEUE,
+ workflow.getWorkflowId(),
+ workflow.getPriority(),
+ properties.getWorkflowOffsetTimeout().getSeconds());
+ executionDAOFacade.updateWorkflow(workflow);
+ decide(workflowId);
}
- public TaskResult fromProto(TaskResultPb.TaskResult from) {
- TaskResult to = new TaskResult();
- to.setWorkflowInstanceId( from.getWorkflowInstanceId() );
- to.setTaskId( from.getTaskId() );
- to.setReasonForIncompletion( from.getReasonForIncompletion() );
- to.setCallbackAfterSeconds( from.getCallbackAfterSeconds() );
- to.setWorkerId( from.getWorkerId() );
- to.setStatus( fromProto( from.getStatus() ) );
- Map outputDataMap = new HashMap();
- for (Map.Entry pair : from.getOutputDataMap().entrySet()) {
- outputDataMap.put( pair.getKey(), fromProto( pair.getValue() ) );
- }
- to.setOutputData(outputDataMap);
- if (from.hasOutputMessage()) {
- to.setOutputMessage( fromProto( from.getOutputMessage() ) );
- }
- return to;
+ /**
+ * @param workflowId the id of the workflow
+ * @param taskReferenceName the referenceName of the task to be skipped
+ * @param skipTaskRequest the {@link SkipTaskRequest} object
+ * @throws IllegalStateException
+ */
+ public void skipTaskFromWorkflow(
+ String workflowId, String taskReferenceName, SkipTaskRequest skipTaskRequest) {
+
+ WorkflowModel workflow = executionDAOFacade.getWorkflowModel(workflowId, true);
+
+ // FIXME Backwards compatibility for legacy workflows already running.
+ // This code will be removed in a future version.
+ workflow = metadataMapperService.populateWorkflowWithDefinitions(workflow);
+
+ // If the workflow is not running then cannot skip any task
+ if (!workflow.getStatus().equals(WorkflowModel.Status.RUNNING)) {
+ String errorMsg =
+ String.format(
+ "The workflow %s is not running so the task referenced by %s cannot be skipped",
+ workflowId, taskReferenceName);
+ throw new IllegalStateException(errorMsg);
+ }
+
+ // Check if the reference name is as per the workflowdef
+ WorkflowTask workflowTask =
+ workflow.getWorkflowDefinition().getTaskByRefName(taskReferenceName);
+ if (workflowTask == null) {
+ String errorMsg =
+ String.format(
+ "The task referenced by %s does not exist in the WorkflowDefinition %s",
+ taskReferenceName, workflow.getWorkflowName());
+ throw new IllegalStateException(errorMsg);
+ }
+
+ // If the task is already started the again it cannot be skipped
+ workflow.getTasks()
+ .forEach(
+ task -> {
+ if (task.getReferenceTaskName().equals(taskReferenceName)) {
+ String errorMsg =
+ String.format(
+ "The task referenced %s has already been processed, cannot be skipped",
+ taskReferenceName);
+ throw new IllegalStateException(errorMsg);
+ }
+ });
+
+ // Now create a "SKIPPED" task for this workflow
+ TaskModel taskToBeSkipped = new TaskModel();
+ taskToBeSkipped.setTaskId(IDGenerator.generate());
+ taskToBeSkipped.setReferenceTaskName(taskReferenceName);
+ taskToBeSkipped.setWorkflowInstanceId(workflowId);
+ taskToBeSkipped.setWorkflowPriority(workflow.getPriority());
+ taskToBeSkipped.setStatus(SKIPPED);
+ taskToBeSkipped.setTaskType(workflowTask.getName());
+ taskToBeSkipped.setCorrelationId(workflow.getCorrelationId());
+ if (skipTaskRequest != null) {
+ taskToBeSkipped.setInputData(skipTaskRequest.getTaskInput());
+ taskToBeSkipped.setOutputData(skipTaskRequest.getTaskOutput());
+ taskToBeSkipped.setInputMessage(skipTaskRequest.getTaskInputMessage());
+ taskToBeSkipped.setOutputMessage(skipTaskRequest.getTaskOutputMessage());
+ }
+ executionDAOFacade.createTasks(Collections.singletonList(taskToBeSkipped));
+ decide(workflowId);
}
- public TaskResultPb.TaskResult.Status toProto(TaskResult.Status from) {
- TaskResultPb.TaskResult.Status to;
- switch (from) {
- case IN_PROGRESS: to = TaskResultPb.TaskResult.Status.IN_PROGRESS; break;
- case FAILED: to = TaskResultPb.TaskResult.Status.FAILED; break;
- case FAILED_WITH_TERMINAL_ERROR: to = TaskResultPb.TaskResult.Status.FAILED_WITH_TERMINAL_ERROR; break;
- case COMPLETED: to = TaskResultPb.TaskResult.Status.COMPLETED; break;
- default: throw new IllegalArgumentException("Unexpected enum constant: " + from);
- }
- return to;
+ public WorkflowModel getWorkflow(String workflowId, boolean includeTasks) {
+ return executionDAOFacade.getWorkflowModel(workflowId, includeTasks);
}
- public TaskResult.Status fromProto(TaskResultPb.TaskResult.Status from) {
- TaskResult.Status to;
- switch (from) {
- case IN_PROGRESS: to = TaskResult.Status.IN_PROGRESS; break;
- case FAILED: to = TaskResult.Status.FAILED; break;
- case FAILED_WITH_TERMINAL_ERROR: to = TaskResult.Status.FAILED_WITH_TERMINAL_ERROR; break;
- case COMPLETED: to = TaskResult.Status.COMPLETED; break;
- default: throw new IllegalArgumentException("Unexpected enum constant: " + from);
- }
- return to;
+ public void addTaskToQueue(TaskModel task) {
+ // put in queue
+ String taskQueueName = QueueUtils.getQueueName(task);
+ if (task.getCallbackAfterSeconds() > 0) {
+ queueDAO.push(
+ taskQueueName,
+ task.getTaskId(),
+ task.getWorkflowPriority(),
+ task.getCallbackAfterSeconds());
+ } else {
+ queueDAO.push(taskQueueName, task.getTaskId(), task.getWorkflowPriority(), 0);
+ }
+ LOGGER.debug(
+ "Added task {} with priority {} to queue {} with call back seconds {}",
+ task,
+ task.getWorkflowPriority(),
+ taskQueueName,
+ task.getCallbackAfterSeconds());
}
- public TaskSummaryPb.TaskSummary toProto(TaskSummary from) {
- TaskSummaryPb.TaskSummary.Builder to = TaskSummaryPb.TaskSummary.newBuilder();
- if (from.getWorkflowId() != null) {
- to.setWorkflowId( from.getWorkflowId() );
- }
- if (from.getWorkflowType() != null) {
- to.setWorkflowType( from.getWorkflowType() );
- }
- if (from.getCorrelationId() != null) {
- to.setCorrelationId( from.getCorrelationId() );
- }
- if (from.getScheduledTime() != null) {
- to.setScheduledTime( from.getScheduledTime() );
- }
- if (from.getStartTime() != null) {
- to.setStartTime( from.getStartTime() );
- }
- if (from.getUpdateTime() != null) {
- to.setUpdateTime( from.getUpdateTime() );
- }
- if (from.getEndTime() != null) {
- to.setEndTime( from.getEndTime() );
- }
- if (from.getStatus() != null) {
- to.setStatus( toProto( from.getStatus() ) );
- }
- if (from.getReasonForIncompletion() != null) {
- to.setReasonForIncompletion( from.getReasonForIncompletion() );
- }
- to.setExecutionTime( from.getExecutionTime() );
- to.setQueueWaitTime( from.getQueueWaitTime() );
- if (from.getTaskDefName() != null) {
- to.setTaskDefName( from.getTaskDefName() );
- }
- if (from.getTaskType() != null) {
- to.setTaskType( from.getTaskType() );
- }
- if (from.getInput() != null) {
- to.setInput( from.getInput() );
+ @VisibleForTesting
+ void setTaskDomains(List tasks, WorkflowModel workflow) {
+ Map taskToDomain = workflow.getTaskToDomain();
+ if (taskToDomain != null) {
+ // Step 1: Apply * mapping to all tasks, if present.
+ String domainstr = taskToDomain.get("*");
+ if (StringUtils.isNotBlank(domainstr)) {
+ String[] domains = domainstr.split(",");
+ tasks.forEach(
+ task -> {
+ // Filter out SystemTask
+ if (!systemTaskRegistry.isSystemTask(task.getTaskType())) {
+ // Check which domain worker is polling
+ // Set the task domain
+ task.setDomain(getActiveDomain(task.getTaskType(), domains));
+ }
+ });
+ }
+ // Step 2: Override additional mappings.
+ tasks.forEach(
+ task -> {
+ if (!systemTaskRegistry.isSystemTask(task.getTaskType())) {
+ String taskDomainstr = taskToDomain.get(task.getTaskType());
+ if (taskDomainstr != null) {
+ task.setDomain(
+ getActiveDomain(
+ task.getTaskType(), taskDomainstr.split(",")));
+ }
+ }
+ });
}
- if (from.getOutput() != null) {
- to.setOutput( from.getOutput() );
- }
- if (from.getTaskId() != null) {
- to.setTaskId( from.getTaskId() );
- }
- if (from.getExternalInputPayloadStoragePath() != null) {
- to.setExternalInputPayloadStoragePath( from.getExternalInputPayloadStoragePath() );
- }
- if (from.getExternalOutputPayloadStoragePath() != null) {
- to.setExternalOutputPayloadStoragePath( from.getExternalOutputPayloadStoragePath() );
- }
- to.setWorkflowPriority( from.getWorkflowPriority() );
- return to.build();
}
- public TaskSummary fromProto(TaskSummaryPb.TaskSummary from) {
- TaskSummary to = new TaskSummary();
- to.setWorkflowId( from.getWorkflowId() );
- to.setWorkflowType( from.getWorkflowType() );
- to.setCorrelationId( from.getCorrelationId() );
- to.setScheduledTime( from.getScheduledTime() );
- to.setStartTime( from.getStartTime() );
- to.setUpdateTime( from.getUpdateTime() );
- to.setEndTime( from.getEndTime() );
- to.setStatus( fromProto( from.getStatus() ) );
- to.setReasonForIncompletion( from.getReasonForIncompletion() );
- to.setExecutionTime( from.getExecutionTime() );
- to.setQueueWaitTime( from.getQueueWaitTime() );
- to.setTaskDefName( from.getTaskDefName() );
- to.setTaskType( from.getTaskType() );
- to.setInput( from.getInput() );
- to.setOutput( from.getOutput() );
- to.setTaskId( from.getTaskId() );
- to.setExternalInputPayloadStoragePath( from.getExternalInputPayloadStoragePath() );
- to.setExternalOutputPayloadStoragePath( from.getExternalOutputPayloadStoragePath() );
- to.setWorkflowPriority( from.getWorkflowPriority() );
- return to;
+ /**
+ * Gets the active domain from the list of domains where the task is to be queued. The domain
+ * list must be ordered. In sequence, check if any worker has polled for last
+ * `activeWorkerLastPollMs`, if so that is the Active domain. When no active domains are found:
+ * If NO_DOMAIN token is provided, return null.
+ * Else, return last domain from list.
+ *
+ * @param taskType the taskType of the task for which active domain is to be found
+ * @param domains the array of domains for the task. (Must contain atleast one element).
+ * @return the active domain where the task will be queued
+ */
+ @VisibleForTesting
+ String getActiveDomain(String taskType, String[] domains) {
+ if (domains == null || domains.length == 0) {
+ return null;
+ }
+
+ return Arrays.stream(domains)
+ .filter(domain -> !domain.equalsIgnoreCase("NO_DOMAIN"))
+ .map(domain -> executionDAOFacade.getTaskPollDataByDomain(taskType, domain.trim()))
+ .filter(Objects::nonNull)
+ .filter(validateLastPolledTime)
+ .findFirst()
+ .map(PollData::getDomain)
+ .orElse(
+ domains[domains.length - 1].trim().equalsIgnoreCase("NO_DOMAIN")
+ ? null
+ : domains[domains.length - 1].trim());
}
- public WorkflowPb.Workflow toProto(Workflow from) {
- WorkflowPb.Workflow.Builder to = WorkflowPb.Workflow.newBuilder();
- if (from.getStatus() != null) {
- to.setStatus( toProto( from.getStatus() ) );
- }
- to.setEndTime( from.getEndTime() );
- if (from.getWorkflowId() != null) {
- to.setWorkflowId( from.getWorkflowId() );
- }
- if (from.getParentWorkflowId() != null) {
- to.setParentWorkflowId( from.getParentWorkflowId() );
- }
- if (from.getParentWorkflowTaskId() != null) {
- to.setParentWorkflowTaskId( from.getParentWorkflowTaskId() );
- }
- for (Task elem : from.getTasks()) {
- to.addTasks( toProto(elem) );
- }
- for (Map.Entry pair : from.getInput().entrySet()) {
- to.putInput( pair.getKey(), toProto( pair.getValue() ) );
- }
- for (Map.Entry pair : from.getOutput().entrySet()) {
- to.putOutput( pair.getKey(), toProto( pair.getValue() ) );
- }
- if (from.getCorrelationId() != null) {
- to.setCorrelationId( from.getCorrelationId() );
- }
- if (from.getReRunFromWorkflowId() != null) {
- to.setReRunFromWorkflowId( from.getReRunFromWorkflowId() );
- }
- if (from.getReasonForIncompletion() != null) {
- to.setReasonForIncompletion( from.getReasonForIncompletion() );
+ private long getTaskDuration(long s, TaskModel task) {
+ long duration = task.getEndTime() - task.getStartTime();
+ s += duration;
+ if (task.getRetriedTaskId() == null) {
+ return s;
}
- if (from.getEvent() != null) {
- to.setEvent( from.getEvent() );
- }
- to.putAllTaskToDomain( from.getTaskToDomain() );
- to.addAllFailedReferenceTaskNames( from.getFailedReferenceTaskNames() );
- if (from.getWorkflowDefinition() != null) {
- to.setWorkflowDefinition( toProto( from.getWorkflowDefinition() ) );
- }
- if (from.getExternalInputPayloadStoragePath() != null) {
- to.setExternalInputPayloadStoragePath( from.getExternalInputPayloadStoragePath() );
- }
- if (from.getExternalOutputPayloadStoragePath() != null) {
- to.setExternalOutputPayloadStoragePath( from.getExternalOutputPayloadStoragePath() );
- }
- to.setPriority( from.getPriority() );
- for (Map.Entry pair : from.getVariables().entrySet()) {
- to.putVariables( pair.getKey(), toProto( pair.getValue() ) );
- }
- to.setLastRetriedTime( from.getLastRetriedTime() );
- return to.build();
- }
-
- public Workflow fromProto(WorkflowPb.Workflow from) {
- Workflow to = new Workflow();
- to.setStatus( fromProto( from.getStatus() ) );
- to.setEndTime( from.getEndTime() );
- to.setWorkflowId( from.getWorkflowId() );
- to.setParentWorkflowId( from.getParentWorkflowId() );
- to.setParentWorkflowTaskId( from.getParentWorkflowTaskId() );
- to.setTasks( from.getTasksList().stream().map(this::fromProto).collect(Collectors.toCollection(ArrayList::new)) );
- Map inputMap = new HashMap();
- for (Map.Entry pair : from.getInputMap().entrySet()) {
- inputMap.put( pair.getKey(), fromProto( pair.getValue() ) );
- }
- to.setInput(inputMap);
- Map outputMap = new HashMap();
- for (Map.Entry pair : from.getOutputMap().entrySet()) {
- outputMap.put( pair.getKey(), fromProto( pair.getValue() ) );
- }
- to.setOutput(outputMap);
- to.setCorrelationId( from.getCorrelationId() );
- to.setReRunFromWorkflowId( from.getReRunFromWorkflowId() );
- to.setReasonForIncompletion( from.getReasonForIncompletion() );
- to.setEvent( from.getEvent() );
- to.setTaskToDomain( from.getTaskToDomainMap() );
- to.setFailedReferenceTaskNames( from.getFailedReferenceTaskNamesList().stream().collect(Collectors.toCollection(HashSet::new)) );
- if (from.hasWorkflowDefinition()) {
- to.setWorkflowDefinition( fromProto( from.getWorkflowDefinition() ) );
- }
- to.setExternalInputPayloadStoragePath( from.getExternalInputPayloadStoragePath() );
- to.setExternalOutputPayloadStoragePath( from.getExternalOutputPayloadStoragePath() );
- to.setPriority( from.getPriority() );
- Map variablesMap = new HashMap();
- for (Map.Entry pair : from.getVariablesMap().entrySet()) {
- variablesMap.put( pair.getKey(), fromProto( pair.getValue() ) );
- }
- to.setVariables(variablesMap);
- to.setLastRetriedTime( from.getLastRetriedTime() );
- return to;
+ return s + getTaskDuration(s, executionDAOFacade.getTaskModel(task.getRetriedTaskId()));
}
- public WorkflowPb.Workflow.WorkflowStatus toProto(Workflow.WorkflowStatus from) {
- WorkflowPb.Workflow.WorkflowStatus to;
- switch (from) {
- case RUNNING: to = WorkflowPb.Workflow.WorkflowStatus.RUNNING; break;
- case COMPLETED: to = WorkflowPb.Workflow.WorkflowStatus.COMPLETED; break;
- case FAILED: to = WorkflowPb.Workflow.WorkflowStatus.FAILED; break;
- case TIMED_OUT: to = WorkflowPb.Workflow.WorkflowStatus.TIMED_OUT; break;
- case TERMINATED: to = WorkflowPb.Workflow.WorkflowStatus.TERMINATED; break;
- case PAUSED: to = WorkflowPb.Workflow.WorkflowStatus.PAUSED; break;
- default: throw new IllegalArgumentException("Unexpected enum constant: " + from);
- }
- return to;
+ @VisibleForTesting
+ boolean scheduleTask(WorkflowModel workflow, List tasks) {
+ List createdTasks;
+ List tasksToBeQueued;
+ boolean startedSystemTasks = false;
+
+ try {
+ if (tasks == null || tasks.isEmpty()) {
+ return false;
+ }
+
+ // Get the highest seq number
+ int count = workflow.getTasks().stream().mapToInt(TaskModel::getSeq).max().orElse(0);
+
+ for (TaskModel task : tasks) {
+ if (task.getSeq() == 0) { // Set only if the seq was not set
+ task.setSeq(++count);
+ }
+ }
+
+ // metric to track the distribution of number of tasks within a workflow
+ Monitors.recordNumTasksInWorkflow(
+ workflow.getTasks().size() + tasks.size(),
+ workflow.getWorkflowName(),
+ String.valueOf(workflow.getWorkflowVersion()));
+
+ // Save the tasks in the DAO
+ createdTasks = executionDAOFacade.createTasks(tasks);
+
+ List systemTasks =
+ createdTasks.stream()
+ .filter(task -> systemTaskRegistry.isSystemTask(task.getTaskType()))
+ .collect(Collectors.toList());
+
+ tasksToBeQueued =
+ createdTasks.stream()
+ .filter(task -> !systemTaskRegistry.isSystemTask(task.getTaskType()))
+ .collect(Collectors.toList());
+
+ // Traverse through all the system tasks, start the sync tasks, in case of async queue
+ // the tasks
+ for (TaskModel task : systemTasks) {
+ WorkflowSystemTask workflowSystemTask = systemTaskRegistry.get(task.getTaskType());
+ if (workflowSystemTask == null) {
+ throw new ApplicationException(
+ NOT_FOUND, "No system task found by name " + task.getTaskType());
+ }
+ if (task.getStatus() != null
+ && !task.getStatus().isTerminal()
+ && task.getStartTime() == 0) {
+ task.setStartTime(System.currentTimeMillis());
+ }
+ if (!workflowSystemTask.isAsync()) {
+ try {
+ workflowSystemTask.start(workflow, task, this);
+ } catch (Exception e) {
+ String errorMsg =
+ String.format(
+ "Unable to start system task: %s, {id: %s, name: %s}",
+ task.getTaskType(),
+ task.getTaskId(),
+ task.getTaskDefName());
+ throw new ApplicationException(
+ ApplicationException.Code.INTERNAL_ERROR, errorMsg, e);
+ }
+ startedSystemTasks = true;
+ executionDAOFacade.updateTask(task);
+ } else {
+ tasksToBeQueued.add(task);
+ }
+ }
+ } catch (Exception e) {
+ List taskIds =
+ tasks.stream().map(TaskModel::getTaskId).collect(Collectors.toList());
+ String errorMsg =
+ String.format(
+ "Error scheduling tasks: %s, for workflow: %s",
+ taskIds, workflow.getWorkflowId());
+ LOGGER.error(errorMsg, e);
+ Monitors.error(CLASS_NAME, "scheduleTask");
+ throw new TerminateWorkflowException(errorMsg);
+ }
+
+ // On addTaskToQueue failures, ignore the exceptions and let WorkflowRepairService take care
+ // of republishing the messages to the queue.
+ try {
+ addTaskToQueue(tasksToBeQueued);
+ } catch (Exception e) {
+ List taskIds =
+ tasksToBeQueued.stream().map(TaskModel::getTaskId).collect(Collectors.toList());
+ String errorMsg =
+ String.format(
+ "Error pushing tasks to the queue: %s, for workflow: %s",
+ taskIds, workflow.getWorkflowId());
+ LOGGER.warn(errorMsg, e);
+ Monitors.error(CLASS_NAME, "scheduleTask");
+ }
+ return startedSystemTasks;
}
- public Workflow.WorkflowStatus fromProto(WorkflowPb.Workflow.WorkflowStatus from) {
- Workflow.WorkflowStatus to;
- switch (from) {
- case RUNNING: to = Workflow.WorkflowStatus.RUNNING; break;
- case COMPLETED: to = Workflow.WorkflowStatus.COMPLETED; break;
- case FAILED: to = Workflow.WorkflowStatus.FAILED; break;
- case TIMED_OUT: to = Workflow.WorkflowStatus.TIMED_OUT; break;
- case TERMINATED: to = Workflow.WorkflowStatus.TERMINATED; break;
- case PAUSED: to = Workflow.WorkflowStatus.PAUSED; break;
- default: throw new IllegalArgumentException("Unexpected enum constant: " + from);
- }
- return to;
+ private void addTaskToQueue(final List tasks) {
+ for (TaskModel task : tasks) {
+ addTaskToQueue(task);
+ }
}
- public WorkflowDefPb.WorkflowDef toProto(WorkflowDef from) {
- WorkflowDefPb.WorkflowDef.Builder to = WorkflowDefPb.WorkflowDef.newBuilder();
- if (from.getName() != null) {
- to.setName( from.getName() );
- }
- if (from.getDescription() != null) {
- to.setDescription( from.getDescription() );
- }
- to.setVersion( from.getVersion() );
- for (WorkflowTask elem : from.getTasks()) {
- to.addTasks( toProto(elem) );
- }
- to.addAllInputParameters( from.getInputParameters() );
- for (Map.Entry pair : from.getOutputParameters().entrySet()) {
- to.putOutputParameters( pair.getKey(), toProto( pair.getValue() ) );
- }
- if (from.getFailureWorkflow() != null) {
- to.setFailureWorkflow( from.getFailureWorkflow() );
+ private WorkflowModel terminate(
+ final WorkflowModel workflow, TerminateWorkflowException terminateWorkflowException) {
+ if (!workflow.getStatus().isTerminal()) {
+ workflow.setStatus(terminateWorkflowException.getWorkflowStatus());
}
- to.setSchemaVersion( from.getSchemaVersion() );
- to.setRestartable( from.isRestartable() );
- to.setWorkflowStatusListenerEnabled( from.isWorkflowStatusListenerEnabled() );
- if (from.getOwnerEmail() != null) {
- to.setOwnerEmail( from.getOwnerEmail() );
- }
- if (from.getTimeoutPolicy() != null) {
- to.setTimeoutPolicy( toProto( from.getTimeoutPolicy() ) );
+
+ if (terminateWorkflowException.getTask() != null && workflow.getFailedTaskId() == null) {
+ workflow.setFailedTaskId(terminateWorkflowException.getTask().getTaskId());
}
- to.setTimeoutSeconds( from.getTimeoutSeconds() );
- for (Map.Entry pair : from.getVariables().entrySet()) {
- to.putVariables( pair.getKey(), toProto( pair.getValue() ) );
+
+ String failureWorkflow = workflow.getWorkflowDefinition().getFailureWorkflow();
+ if (failureWorkflow != null) {
+ if (failureWorkflow.startsWith("$")) {
+ String[] paramPathComponents = failureWorkflow.split("\\.");
+ String name = paramPathComponents[2]; // name of the input parameter
+ failureWorkflow = (String) workflow.getInput().get(name);
+ }
}
- for (Map.Entry pair : from.getInputTemplate().entrySet()) {
- to.putInputTemplate( pair.getKey(), toProto( pair.getValue() ) );
+ if (terminateWorkflowException.getTask() != null) {
+ executionDAOFacade.updateTask(terminateWorkflowException.getTask());
}
- return to.build();
+ return terminateWorkflow(
+ workflow, terminateWorkflowException.getMessage(), failureWorkflow);
}
- public WorkflowDef fromProto(WorkflowDefPb.WorkflowDef from) {
- WorkflowDef to = new WorkflowDef();
- to.setName( from.getName() );
- to.setDescription( from.getDescription() );
- to.setVersion( from.getVersion() );
- to.setTasks( from.getTasksList().stream().map(this::fromProto).collect(Collectors.toCollection(ArrayList::new)) );
- to.setInputParameters( from.getInputParametersList().stream().collect(Collectors.toCollection(ArrayList::new)) );
- Map outputParametersMap = new HashMap();
- for (Map.Entry pair : from.getOutputParametersMap().entrySet()) {
- outputParametersMap.put( pair.getKey(), fromProto( pair.getValue() ) );
- }
- to.setOutputParameters(outputParametersMap);
- to.setFailureWorkflow( from.getFailureWorkflow() );
- to.setSchemaVersion( from.getSchemaVersion() );
- to.setRestartable( from.getRestartable() );
- to.setWorkflowStatusListenerEnabled( from.getWorkflowStatusListenerEnabled() );
- to.setOwnerEmail( from.getOwnerEmail() );
- to.setTimeoutPolicy( fromProto( from.getTimeoutPolicy() ) );
- to.setTimeoutSeconds( from.getTimeoutSeconds() );
- Map variablesMap = new HashMap();
- for (Map.Entry pair : from.getVariablesMap().entrySet()) {
- variablesMap.put( pair.getKey(), fromProto( pair.getValue() ) );
- }
- to.setVariables(variablesMap);
- Map inputTemplateMap = new HashMap();
- for (Map.Entry pair : from.getInputTemplateMap().entrySet()) {
- inputTemplateMap.put( pair.getKey(), fromProto( pair.getValue() ) );
- }
- to.setInputTemplate(inputTemplateMap);
- return to;
+ private boolean rerunWF(
+ String workflowId,
+ String taskId,
+ Map taskInput,
+ Map workflowInput,
+ String correlationId) {
+
+ // Get the workflow
+ WorkflowModel workflow = executionDAOFacade.getWorkflowModel(workflowId, true);
+ updateAndPushParents(workflow, "reran");
+
+ // If the task Id is null it implies that the entire workflow has to be rerun
+ if (taskId == null) {
+ // remove all tasks
+ workflow.getTasks().forEach(task -> executionDAOFacade.removeTask(task.getTaskId()));
+ // Set workflow as RUNNING
+ workflow.setStatus(WorkflowModel.Status.RUNNING);
+ // Reset failure reason from previous run to default
+ workflow.setReasonForIncompletion(null);
+ workflow.setFailedTaskId(null);
+ workflow.setFailedReferenceTaskNames(new HashSet<>());
+ if (correlationId != null) {
+ workflow.setCorrelationId(correlationId);
+ }
+ if (workflowInput != null) {
+ workflow.setInput(workflowInput);
+ }
+
+ queueDAO.push(
+ DECIDER_QUEUE,
+ workflow.getWorkflowId(),
+ workflow.getPriority(),
+ properties.getWorkflowOffsetTimeout().getSeconds());
+ executionDAOFacade.updateWorkflow(workflow);
+
+ decide(workflowId);
+ return true;
+ }
+
+ // Now iterate through the tasks and find the "specific" task
+ TaskModel rerunFromTask = null;
+ for (TaskModel task : workflow.getTasks()) {
+ if (task.getTaskId().equals(taskId)) {
+ rerunFromTask = task;
+ break;
+ }
+ }
+
+ // If not found look into sub workflows
+ if (rerunFromTask == null) {
+ for (TaskModel task : workflow.getTasks()) {
+ if (task.getTaskType().equalsIgnoreCase(TaskType.TASK_TYPE_SUB_WORKFLOW)) {
+ String subWorkflowId = task.getSubWorkflowId();
+ if (rerunWF(subWorkflowId, taskId, taskInput, null, null)) {
+ rerunFromTask = task;
+ break;
+ }
+ }
+ }
+ }
+
+ if (rerunFromTask != null) {
+ // set workflow as RUNNING
+ workflow.setStatus(WorkflowModel.Status.RUNNING);
+ // Reset failure reason from previous run to default
+ workflow.setReasonForIncompletion(null);
+ workflow.setFailedTaskId(null);
+ workflow.setFailedReferenceTaskNames(new HashSet<>());
+ if (correlationId != null) {
+ workflow.setCorrelationId(correlationId);
+ }
+ if (workflowInput != null) {
+ workflow.setInput(workflowInput);
+ }
+ // Add to decider queue
+ queueDAO.push(
+ DECIDER_QUEUE,
+ workflow.getWorkflowId(),
+ workflow.getPriority(),
+ properties.getWorkflowOffsetTimeout().getSeconds());
+ executionDAOFacade.updateWorkflow(workflow);
+ // update tasks in datastore to update workflow-tasks relationship for archived
+ // workflows
+ executionDAOFacade.updateTasks(workflow.getTasks());
+ // Remove all tasks after the "rerunFromTask"
+ for (TaskModel task : workflow.getTasks()) {
+ if (task.getSeq() > rerunFromTask.getSeq()) {
+ executionDAOFacade.removeTask(task.getTaskId());
+ }
+ }
+ // reset fields before restarting the task
+ rerunFromTask.setScheduledTime(System.currentTimeMillis());
+ rerunFromTask.setStartTime(0);
+ rerunFromTask.setUpdateTime(0);
+ rerunFromTask.setEndTime(0);
+ rerunFromTask.getOutputData().clear();
+ rerunFromTask.setRetried(false);
+ rerunFromTask.setExecuted(false);
+ rerunFromTask.setExternalOutputPayloadStoragePath(null);
+ if (rerunFromTask.getTaskType().equalsIgnoreCase(TaskType.TASK_TYPE_SUB_WORKFLOW)) {
+ // if task is sub workflow set task as IN_PROGRESS and reset start time
+ rerunFromTask.setStatus(IN_PROGRESS);
+ rerunFromTask.setStartTime(System.currentTimeMillis());
+ } else {
+ if (taskInput != null) {
+ rerunFromTask.setInputData(taskInput);
+ }
+ if (systemTaskRegistry.isSystemTask(rerunFromTask.getTaskType())
+ && !systemTaskRegistry.get(rerunFromTask.getTaskType()).isAsync()) {
+ // Start the synchronous system task directly
+ systemTaskRegistry
+ .get(rerunFromTask.getTaskType())
+ .start(workflow, rerunFromTask, this);
+ } else {
+ // Set the task to rerun as SCHEDULED
+ rerunFromTask.setStatus(SCHEDULED);
+ addTaskToQueue(rerunFromTask);
+ }
+ }
+ executionDAOFacade.updateTask(rerunFromTask);
+
+ decide(workflowId);
+ return true;
+ }
+ return false;
}
- public WorkflowDefPb.WorkflowDef.TimeoutPolicy toProto(WorkflowDef.TimeoutPolicy from) {
- WorkflowDefPb.WorkflowDef.TimeoutPolicy to;
- switch (from) {
- case TIME_OUT_WF: to = WorkflowDefPb.WorkflowDef.TimeoutPolicy.TIME_OUT_WF; break;
- case ALERT_ONLY: to = WorkflowDefPb.WorkflowDef.TimeoutPolicy.ALERT_ONLY; break;
- default: throw new IllegalArgumentException("Unexpected enum constant: " + from);
- }
- return to;
+ public void scheduleNextIteration(TaskModel loopTask, WorkflowModel workflow) {
+ // Schedule only first loop over task. Rest will be taken care in Decider Service when this
+ // task will get completed.
+ List scheduledLoopOverTasks =
+ deciderService.getTasksToBeScheduled(
+ workflow,
+ loopTask.getWorkflowTask().getLoopOver().get(0),
+ loopTask.getRetryCount(),
+ null);
+ setTaskDomains(scheduledLoopOverTasks, workflow);
+ scheduledLoopOverTasks.forEach(
+ t -> {
+ t.setReferenceTaskName(
+ TaskUtils.appendIteration(
+ t.getReferenceTaskName(), loopTask.getIteration()));
+ t.setIteration(loopTask.getIteration());
+ });
+ scheduleTask(workflow, scheduledLoopOverTasks);
}
- public WorkflowDef.TimeoutPolicy fromProto(WorkflowDefPb.WorkflowDef.TimeoutPolicy from) {
- WorkflowDef.TimeoutPolicy to;
- switch (from) {
- case TIME_OUT_WF: to = WorkflowDef.TimeoutPolicy.TIME_OUT_WF; break;
- case ALERT_ONLY: to = WorkflowDef.TimeoutPolicy.ALERT_ONLY; break;
- default: throw new IllegalArgumentException("Unexpected enum constant: " + from);
- }
- return to;
+ public TaskDef getTaskDefinition(TaskModel task) {
+ return task.getTaskDefinition()
+ .orElseGet(
+ () ->
+ Optional.ofNullable(
+ metadataDAO.getTaskDef(
+ task.getWorkflowTask().getName()))
+ .orElseThrow(
+ () -> {
+ String reason =
+ String.format(
+ "Invalid task specified. Cannot find task by name %s in the task definitions",
+ task.getWorkflowTask()
+ .getName());
+ return new TerminateWorkflowException(reason);
+ }));
}
- public WorkflowSummaryPb.WorkflowSummary toProto(WorkflowSummary from) {
- WorkflowSummaryPb.WorkflowSummary.Builder to = WorkflowSummaryPb.WorkflowSummary.newBuilder();
- if (from.getWorkflowType() != null) {
- to.setWorkflowType( from.getWorkflowType() );
- }
- to.setVersion( from.getVersion() );
- if (from.getWorkflowId() != null) {
- to.setWorkflowId( from.getWorkflowId() );
- }
- if (from.getCorrelationId() != null) {
- to.setCorrelationId( from.getCorrelationId() );
- }
- if (from.getStartTime() != null) {
- to.setStartTime( from.getStartTime() );
- }
- if (from.getUpdateTime() != null) {
- to.setUpdateTime( from.getUpdateTime() );
- }
- if (from.getEndTime() != null) {
- to.setEndTime( from.getEndTime() );
- }
- if (from.getStatus() != null) {
- to.setStatus( toProto( from.getStatus() ) );
- }
- if (from.getInput() != null) {
- to.setInput( from.getInput() );
- }
- if (from.getOutput() != null) {
- to.setOutput( from.getOutput() );
- }
- if (from.getReasonForIncompletion() != null) {
- to.setReasonForIncompletion( from.getReasonForIncompletion() );
- }
- to.setExecutionTime( from.getExecutionTime() );
- if (from.getEvent() != null) {
- to.setEvent( from.getEvent() );
- }
- if (from.getFailedReferenceTaskNames() != null) {
- to.setFailedReferenceTaskNames( from.getFailedReferenceTaskNames() );
- }
- if (from.getExternalInputPayloadStoragePath() != null) {
- to.setExternalInputPayloadStoragePath( from.getExternalInputPayloadStoragePath() );
- }
- if (from.getExternalOutputPayloadStoragePath() != null) {
- to.setExternalOutputPayloadStoragePath( from.getExternalOutputPayloadStoragePath() );
- }
- to.setPriority( from.getPriority() );
- return to.build();
+ @VisibleForTesting
+ void updateParentWorkflowTask(WorkflowModel subWorkflow) {
+ TaskModel subWorkflowTask =
+ executionDAOFacade.getTaskModel(subWorkflow.getParentWorkflowTaskId());
+ executeSubworkflowTaskAndSyncData(subWorkflow, subWorkflowTask);
+ executionDAOFacade.updateTask(subWorkflowTask);
}
- public WorkflowSummary fromProto(WorkflowSummaryPb.WorkflowSummary from) {
- WorkflowSummary to = new WorkflowSummary();
- to.setWorkflowType( from.getWorkflowType() );
- to.setVersion( from.getVersion() );
- to.setWorkflowId( from.getWorkflowId() );
- to.setCorrelationId( from.getCorrelationId() );
- to.setStartTime( from.getStartTime() );
- to.setUpdateTime( from.getUpdateTime() );
- to.setEndTime( from.getEndTime() );
- to.setStatus( fromProto( from.getStatus() ) );
- to.setInput( from.getInput() );
- to.setOutput( from.getOutput() );
- to.setReasonForIncompletion( from.getReasonForIncompletion() );
- to.setExecutionTime( from.getExecutionTime() );
- to.setEvent( from.getEvent() );
- to.setFailedReferenceTaskNames( from.getFailedReferenceTaskNames() );
- to.setExternalInputPayloadStoragePath( from.getExternalInputPayloadStoragePath() );
- to.setExternalOutputPayloadStoragePath( from.getExternalOutputPayloadStoragePath() );
- to.setPriority( from.getPriority() );
- return to;
+ private void executeSubworkflowTaskAndSyncData(
+ WorkflowModel subWorkflow, TaskModel subWorkflowTask) {
+ WorkflowSystemTask subWorkflowSystemTask =
+ systemTaskRegistry.get(TaskType.TASK_TYPE_SUB_WORKFLOW);
+ subWorkflowSystemTask.execute(subWorkflow, subWorkflowTask, this);
}
- public WorkflowTaskPb.WorkflowTask toProto(WorkflowTask from) {
- WorkflowTaskPb.WorkflowTask.Builder to = WorkflowTaskPb.WorkflowTask.newBuilder();
- if (from.getName() != null) {
- to.setName( from.getName() );
- }
- if (from.getTaskReferenceName() != null) {
- to.setTaskReferenceName( from.getTaskReferenceName() );
- }
- if (from.getDescription() != null) {
- to.setDescription( from.getDescription() );
- }
- for (Map.Entry pair : from.getInputParameters().entrySet()) {
- to.putInputParameters( pair.getKey(), toProto( pair.getValue() ) );
- }
- if (from.getType() != null) {
- to.setType( from.getType() );
- }
- if (from.getDynamicTaskNameParam() != null) {
- to.setDynamicTaskNameParam( from.getDynamicTaskNameParam() );
- }
- if (from.getCaseValueParam() != null) {
- to.setCaseValueParam( from.getCaseValueParam() );
+ /** Pushes parent workflow id into the decider queue with a priority. */
+ private void pushParentWorkflow(String parentWorkflowId) {
+ if (queueDAO.containsMessage(DECIDER_QUEUE, parentWorkflowId)) {
+ queueDAO.postpone(DECIDER_QUEUE, parentWorkflowId, PARENT_WF_PRIORITY, 0);
+ } else {
+ queueDAO.push(DECIDER_QUEUE, parentWorkflowId, PARENT_WF_PRIORITY, 0);
}
- if (from.getCaseExpression() != null) {
- to.setCaseExpression( from.getCaseExpression() );
- }
- if (from.getScriptExpression() != null) {
- to.setScriptExpression( from.getScriptExpression() );
- }
- for (Map.Entry> pair : from.getDecisionCases().entrySet()) {
- to.putDecisionCases( pair.getKey(), toProto( pair.getValue() ) );
- }
- if (from.getDynamicForkTasksParam() != null) {
- to.setDynamicForkTasksParam( from.getDynamicForkTasksParam() );
- }
- if (from.getDynamicForkTasksInputParamName() != null) {
- to.setDynamicForkTasksInputParamName( from.getDynamicForkTasksInputParamName() );
- }
- for (WorkflowTask elem : from.getDefaultCase()) {
- to.addDefaultCase( toProto(elem) );
- }
- for (List elem : from.getForkTasks()) {
- to.addForkTasks( toProto(elem) );
- }
- to.setStartDelay( from.getStartDelay() );
- if (from.getSubWorkflowParam() != null) {
- to.setSubWorkflowParam( toProto( from.getSubWorkflowParam() ) );
- }
- to.addAllJoinOn( from.getJoinOn() );
- if (from.getSink() != null) {
- to.setSink( from.getSink() );
- }
- to.setOptional( from.isOptional() );
- if (from.getTaskDefinition() != null) {
- to.setTaskDefinition( toProto( from.getTaskDefinition() ) );
- }
- if (from.isRateLimited() != null) {
- to.setRateLimited( from.isRateLimited() );
- }
- to.addAllDefaultExclusiveJoinTask( from.getDefaultExclusiveJoinTask() );
- if (from.isAsyncComplete() != null) {
- to.setAsyncComplete( from.isAsyncComplete() );
- }
- if (from.getLoopCondition() != null) {
- to.setLoopCondition( from.getLoopCondition() );
- }
- for (WorkflowTask elem : from.getLoopOver()) {
- to.addLoopOver( toProto(elem) );
- }
- if (from.getRetryCount() != null) {
- to.setRetryCount( from.getRetryCount() );
- }
- if (from.getEvaluatorType() != null) {
- to.setEvaluatorType( from.getEvaluatorType() );
- }
- if (from.getExpression() != null) {
- to.setExpression( from.getExpression() );
- }
- return to.build();
- }
- public WorkflowTask fromProto(WorkflowTaskPb.WorkflowTask from) {
- WorkflowTask to = new WorkflowTask();
- to.setName( from.getName() );
- to.setTaskReferenceName( from.getTaskReferenceName() );
- to.setDescription( from.getDescription() );
- Map inputParametersMap = new HashMap();
- for (Map.Entry pair : from.getInputParametersMap().entrySet()) {
- inputParametersMap.put( pair.getKey(), fromProto( pair.getValue() ) );
- }
- to.setInputParameters(inputParametersMap);
- to.setType( from.getType() );
- to.setDynamicTaskNameParam( from.getDynamicTaskNameParam() );
- to.setCaseValueParam( from.getCaseValueParam() );
- to.setCaseExpression( from.getCaseExpression() );
- to.setScriptExpression( from.getScriptExpression() );
- Map> decisionCasesMap = new HashMap>();
- for (Map.Entry pair : from.getDecisionCasesMap().entrySet()) {
- decisionCasesMap.put( pair.getKey(), fromProto( pair.getValue() ) );
- }
- to.setDecisionCases(decisionCasesMap);
- to.setDynamicForkTasksParam( from.getDynamicForkTasksParam() );
- to.setDynamicForkTasksInputParamName( from.getDynamicForkTasksInputParamName() );
- to.setDefaultCase( from.getDefaultCaseList().stream().map(this::fromProto).collect(Collectors.toCollection(ArrayList::new)) );
- to.setForkTasks( from.getForkTasksList().stream().map(this::fromProto).collect(Collectors.toCollection(ArrayList::new)) );
- to.setStartDelay( from.getStartDelay() );
- if (from.hasSubWorkflowParam()) {
- to.setSubWorkflowParam( fromProto( from.getSubWorkflowParam() ) );
- }
- to.setJoinOn( from.getJoinOnList().stream().collect(Collectors.toCollection(ArrayList::new)) );
- to.setSink( from.getSink() );
- to.setOptional( from.getOptional() );
- if (from.hasTaskDefinition()) {
- to.setTaskDefinition( fromProto( from.getTaskDefinition() ) );
- }
- to.setRateLimited( from.getRateLimited() );
- to.setDefaultExclusiveJoinTask( from.getDefaultExclusiveJoinTaskList().stream().collect(Collectors.toCollection(ArrayList::new)) );
- to.setAsyncComplete( from.getAsyncComplete() );
- to.setLoopCondition( from.getLoopCondition() );
- to.setLoopOver( from.getLoopOverList().stream().map(this::fromProto).collect(Collectors.toCollection(ArrayList::new)) );
- to.setRetryCount( from.getRetryCount() );
- to.setEvaluatorType( from.getEvaluatorType() );
- to.setExpression( from.getExpression() );
- return to;
+ LOGGER.info("Pushed parent workflow {} to {}", parentWorkflowId, DECIDER_QUEUE);
}
-
- public abstract WorkflowTaskPb.WorkflowTask.WorkflowTaskList toProto(List in);
-
- public abstract List fromProto(WorkflowTaskPb.WorkflowTask.WorkflowTaskList in);
-
- public abstract Value toProto(Object in);
-
- public abstract Object fromProto(Value in);
-
- public abstract Any toProto(Any in);
-
- public abstract Any fromProto(Any in);
}