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: + *

+ */ + 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); }