diff --git a/.idea/compiler.xml b/.idea/compiler.xml index f34f517..7c3c166 100644 --- a/.idea/compiler.xml +++ b/.idea/compiler.xml @@ -15,6 +15,7 @@ + diff --git a/database/migrations/workflow/V2__add_branch_type_to_workflow_edges.sql b/database/migrations/workflow/V2__add_branch_type_to_workflow_edges.sql new file mode 100644 index 0000000..342e438 --- /dev/null +++ b/database/migrations/workflow/V2__add_branch_type_to_workflow_edges.sql @@ -0,0 +1,2 @@ +ALTER TABLE workflow_edges + ADD COLUMN branch_type VARCHAR(20) NOT NULL DEFAULT 'DEFAULT'; \ No newline at end of file diff --git a/database/migrations/workflow/V3__add_output_to_task_executions.sql b/database/migrations/workflow/V3__add_output_to_task_executions.sql new file mode 100644 index 0000000..e0eb8e9 --- /dev/null +++ b/database/migrations/workflow/V3__add_output_to_task_executions.sql @@ -0,0 +1,2 @@ +ALTER TABLE task_executions + ADD COLUMN output jsonb; \ No newline at end of file diff --git a/docker-compose.yml b/docker-compose.yml index f2a0ac7..63c1b25 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -82,8 +82,6 @@ services: SLACK_BOT_TOKEN : ${SLACK_BOT_TOKEN} volumes: - ./database/migrations:/app/database/migrations - environment: - SLACK_BOT_TOKEN: ${SLACK_BOT_TOKEN} kafka: image: apache/kafka:latest diff --git a/workflow-service/src/main/java/com/flowforge/workflowservice/application/execution/condtion/ConditionEvaluator.java b/workflow-service/src/main/java/com/flowforge/workflowservice/application/execution/condtion/ConditionEvaluator.java new file mode 100644 index 0000000..10eaf16 --- /dev/null +++ b/workflow-service/src/main/java/com/flowforge/workflowservice/application/execution/condtion/ConditionEvaluator.java @@ -0,0 +1,53 @@ +package com.flowforge.workflowservice.application.execution.condtion; + +import com.fasterxml.jackson.databind.JsonNode; +import com.flowforge.workflowservice.common.exception.BusinessException; +import com.flowforge.workflowservice.common.exception.ErrorCode; +import com.flowforge.workflowservice.presentation.dto.request.ConditionConfig; +import org.springframework.stereotype.Component; + +@Component +public class ConditionEvaluator { + + public boolean evaluate( + JsonNode output, + ConditionConfig config + ) { + String field = config.field(); + ConditionOperator operator = config.operator(); + JsonNode expectedValue = config.value(); + + JsonNode actualValue = output.get(field); + + if (actualValue == null || actualValue.isNull()){ + throw new BusinessException(ErrorCode.CONDITION_FIELD_NOT_FOUND); + } + + return switch(operator) { + case EQUALS -> actualValue.equals(expectedValue); + case NOT_EQUALS -> !actualValue.equals(expectedValue); + case GREATER_THAN -> { + validateNumeric(actualValue, expectedValue); + yield actualValue.asDouble() > expectedValue.asDouble(); + } + case LESS_THAN -> { + validateNumeric(actualValue, expectedValue); + yield actualValue.asDouble() < expectedValue.asDouble(); + } + case GREATER_OR_EQUAL -> { + validateNumeric(actualValue, expectedValue); + yield actualValue.asDouble() >= expectedValue.asDouble(); + } + case LESS_OR_EQUAL -> { + validateNumeric(actualValue, expectedValue); + yield actualValue.asDouble() <= expectedValue.asDouble(); + } + }; + } + + private void validateNumeric(JsonNode actual, JsonNode expected) { + if (!actual.isNumber() || !expected.isNumber()) { + throw new BusinessException(ErrorCode.INVALID_CONDITION_COMPARISON); + } + } +} \ No newline at end of file diff --git a/workflow-service/src/main/java/com/flowforge/workflowservice/application/execution/condtion/ConditionExecutionService.java b/workflow-service/src/main/java/com/flowforge/workflowservice/application/execution/condtion/ConditionExecutionService.java new file mode 100644 index 0000000..97f952c --- /dev/null +++ b/workflow-service/src/main/java/com/flowforge/workflowservice/application/execution/condtion/ConditionExecutionService.java @@ -0,0 +1,31 @@ +package com.flowforge.workflowservice.application.execution.condtion; + +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; +import com.flowforge.workflowservice.domain.edge.BranchType; +import com.flowforge.workflowservice.domain.node.WorkflowNode; +import com.flowforge.workflowservice.presentation.dto.request.ConditionConfig; +import lombok.RequiredArgsConstructor; +import org.springframework.stereotype.Service; + +@Service +@RequiredArgsConstructor +public class ConditionExecutionService { + private final ObjectMapper objectMapper; + private final ConditionEvaluator conditionEvaluator; + + public BranchType execute( + JsonNode previousTaskOutput, + WorkflowNode conditionNode + ){ + ConditionConfig config = objectMapper.convertValue( + conditionNode.getConfiguration(), + ConditionConfig.class + ); + + boolean result = conditionEvaluator.evaluate(previousTaskOutput,config); + return result + ? BranchType.TRUE + : BranchType.FALSE; + } +} diff --git a/workflow-service/src/main/java/com/flowforge/workflowservice/application/execution/condtion/ConditionOperator.java b/workflow-service/src/main/java/com/flowforge/workflowservice/application/execution/condtion/ConditionOperator.java new file mode 100644 index 0000000..14e326b --- /dev/null +++ b/workflow-service/src/main/java/com/flowforge/workflowservice/application/execution/condtion/ConditionOperator.java @@ -0,0 +1,10 @@ +package com.flowforge.workflowservice.application.execution.condtion; + +public enum ConditionOperator { + EQUALS, + NOT_EQUALS, + GREATER_THAN, + LESS_THAN, + GREATER_OR_EQUAL, + LESS_OR_EQUAL +} diff --git a/workflow-service/src/main/java/com/flowforge/workflowservice/application/execution/engine/ExecutionEngine.java b/workflow-service/src/main/java/com/flowforge/workflowservice/application/execution/engine/ExecutionEngine.java index 932f339..7b49161 100644 --- a/workflow-service/src/main/java/com/flowforge/workflowservice/application/execution/engine/ExecutionEngine.java +++ b/workflow-service/src/main/java/com/flowforge/workflowservice/application/execution/engine/ExecutionEngine.java @@ -39,8 +39,7 @@ public void handleTaskCompleted(TaskExecution taskExecution) { ); List runnableNodes = dependencyResolver.resolveNextNodes( - taskExecution.getWorkflowExecution(), - taskExecution.getNode()); + taskExecution); log.info( "Resolved next nodes {}", diff --git a/workflow-service/src/main/java/com/flowforge/workflowservice/application/execution/resolver/DefaultDependencyResolver.java b/workflow-service/src/main/java/com/flowforge/workflowservice/application/execution/resolver/DefaultDependencyResolver.java index abc3056..11d044d 100644 --- a/workflow-service/src/main/java/com/flowforge/workflowservice/application/execution/resolver/DefaultDependencyResolver.java +++ b/workflow-service/src/main/java/com/flowforge/workflowservice/application/execution/resolver/DefaultDependencyResolver.java @@ -1,13 +1,18 @@ package com.flowforge.workflowservice.application.execution.resolver; -import com.flowforge.workflowservice.application.execution.scheduler.NodeScheduler; +import com.flowforge.workflowservice.application.execution.condtion.ConditionExecutionService; +import com.flowforge.workflowservice.common.exception.BusinessException; +import com.flowforge.workflowservice.common.exception.ErrorCode; +import com.flowforge.workflowservice.domain.edge.BranchType; import com.flowforge.workflowservice.domain.edge.WorkflowEdge; -import com.flowforge.workflowservice.domain.execution.WorkflowExecution; +import com.flowforge.workflowservice.domain.execution.TaskExecution; +import com.flowforge.workflowservice.domain.node.NodeType; import com.flowforge.workflowservice.domain.node.WorkflowNode; import com.flowforge.workflowservice.infrastructure.persistence.WorkflowEdgeRepository; import lombok.RequiredArgsConstructor; import org.springframework.stereotype.Service; +import java.util.ArrayList; import java.util.List; @Service @@ -15,15 +20,36 @@ public class DefaultDependencyResolver implements DependencyResolver { private final WorkflowEdgeRepository workflowEdgeRepository; + private final ConditionExecutionService conditionExecutionService; @Override public List resolveNextNodes( - WorkflowExecution execution, - WorkflowNode completedNode + TaskExecution completedTask ) { - List edges = workflowEdgeRepository.findBySourceNode(completedNode); - return edges.stream() - .map(WorkflowEdge::getTargetNode) - .toList(); + WorkflowNode completedNode = completedTask.getNode(); + + List nextNodes = new ArrayList<>(); + + List edges = workflowEdgeRepository.findBySourceNode(completedNode); + + for (WorkflowEdge edge : edges){ + WorkflowNode targetNode = edge.getTargetNode(); + if (targetNode.getNodeType() != NodeType.CONDITION){ + nextNodes.add(targetNode); + continue; + } + BranchType branch = conditionExecutionService.execute(completedTask.getOutput(),targetNode); + + WorkflowEdge selectedEdge = workflowEdgeRepository + .findBySourceNodeAndBranchType(targetNode,branch) + .orElseThrow(() -> + new BusinessException( + ErrorCode.CONDITION_BRANCH_NOT_FOUND + )); + + nextNodes.add(selectedEdge.getTargetNode()); + + } + return nextNodes; } } diff --git a/workflow-service/src/main/java/com/flowforge/workflowservice/application/execution/resolver/DependencyResolver.java b/workflow-service/src/main/java/com/flowforge/workflowservice/application/execution/resolver/DependencyResolver.java index 8f9d820..5118bac 100644 --- a/workflow-service/src/main/java/com/flowforge/workflowservice/application/execution/resolver/DependencyResolver.java +++ b/workflow-service/src/main/java/com/flowforge/workflowservice/application/execution/resolver/DependencyResolver.java @@ -1,5 +1,6 @@ package com.flowforge.workflowservice.application.execution.resolver; +import com.flowforge.workflowservice.domain.execution.TaskExecution; import com.flowforge.workflowservice.domain.execution.WorkflowExecution; import com.flowforge.workflowservice.domain.node.WorkflowNode; @@ -8,8 +9,7 @@ public interface DependencyResolver { List resolveNextNodes( - WorkflowExecution execution, - WorkflowNode completedNode + TaskExecution taskExecution ); } \ No newline at end of file diff --git a/workflow-service/src/main/java/com/flowforge/workflowservice/application/execution/scheduler/NodeScheduler.java b/workflow-service/src/main/java/com/flowforge/workflowservice/application/execution/scheduler/NodeScheduler.java index 846e88f..a2c49a6 100644 --- a/workflow-service/src/main/java/com/flowforge/workflowservice/application/execution/scheduler/NodeScheduler.java +++ b/workflow-service/src/main/java/com/flowforge/workflowservice/application/execution/scheduler/NodeScheduler.java @@ -7,6 +7,7 @@ import com.flowforge.workflowservice.common.exception.ErrorCode; import com.flowforge.workflowservice.domain.execution.TaskExecution; import com.flowforge.workflowservice.domain.execution.WorkflowExecution; +import com.flowforge.workflowservice.domain.node.NodeType; import com.flowforge.workflowservice.domain.node.WorkflowNode; import com.flowforge.workflowservice.infrastructure.persistence.TaskExecutionRepository; import lombok.RequiredArgsConstructor; @@ -30,6 +31,7 @@ public void schedule( WorkflowNode node ) { log.info("Entered NodeScheduler"); + TaskExecution taskExecution = taskExecutionRepository .findByWorkflowExecutionIdAndNodeId( @@ -44,6 +46,7 @@ public void schedule( node.getNodeType() ); + taskStateMachine.startTaskExecution(taskExecution); taskExecutionRepository.save(taskExecution); diff --git a/workflow-service/src/main/java/com/flowforge/workflowservice/common/exception/ErrorCode.java b/workflow-service/src/main/java/com/flowforge/workflowservice/common/exception/ErrorCode.java index 862bd54..9043587 100644 --- a/workflow-service/src/main/java/com/flowforge/workflowservice/common/exception/ErrorCode.java +++ b/workflow-service/src/main/java/com/flowforge/workflowservice/common/exception/ErrorCode.java @@ -26,7 +26,11 @@ public enum ErrorCode { INVALID_WORKFLOW_STATE_TRANSITION(HttpStatus.BAD_REQUEST, "Invalid workflow execution state transition"), INVALID_TASK_STATE_TRANSITION(HttpStatus.BAD_REQUEST, "Invalid task execution state transition"), TASK_EXECUTION_NOT_FOUND(HttpStatus.NOT_FOUND, "Task execution not found"), - INVALID_EMAIL_CONFIGURATION(HttpStatus.BAD_REQUEST,"Invalid Email Configuration" ); + INVALID_EMAIL_CONFIGURATION(HttpStatus.BAD_REQUEST,"Invalid Email Configuration" ), + CONDITION_FIELD_NOT_FOUND(HttpStatus.NOT_FOUND,"Condition field not found in task output"), + CONDITION_BRANCH_NOT_FOUND(HttpStatus.NOT_FOUND,"condition branch not found"), + INVALID_CONDITION_COMPARISON(HttpStatus.BAD_REQUEST, "Invalid condition comparison"), + INVALID_CONDITION_OPERATOR(HttpStatus.BAD_REQUEST, "Unsupported condition operator"); private final HttpStatus status; private final String message; diff --git a/workflow-service/src/main/java/com/flowforge/workflowservice/domain/edge/BranchType.java b/workflow-service/src/main/java/com/flowforge/workflowservice/domain/edge/BranchType.java new file mode 100644 index 0000000..0bfaf56 --- /dev/null +++ b/workflow-service/src/main/java/com/flowforge/workflowservice/domain/edge/BranchType.java @@ -0,0 +1,7 @@ +package com.flowforge.workflowservice.domain.edge; + +public enum BranchType { + TRUE, + FALSE, + DEFAULT +} diff --git a/workflow-service/src/main/java/com/flowforge/workflowservice/domain/edge/WorkflowEdge.java b/workflow-service/src/main/java/com/flowforge/workflowservice/domain/edge/WorkflowEdge.java index c3dcb87..a75b621 100644 --- a/workflow-service/src/main/java/com/flowforge/workflowservice/domain/edge/WorkflowEdge.java +++ b/workflow-service/src/main/java/com/flowforge/workflowservice/domain/edge/WorkflowEdge.java @@ -33,4 +33,10 @@ public class WorkflowEdge { @CreationTimestamp private Instant createdAt; + + @Builder.Default + @Enumerated(EnumType.STRING) + @Column(name = "branch_type", nullable = false) + private BranchType branchType = BranchType.DEFAULT; + } \ No newline at end of file diff --git a/workflow-service/src/main/java/com/flowforge/workflowservice/domain/execution/TaskExecution.java b/workflow-service/src/main/java/com/flowforge/workflowservice/domain/execution/TaskExecution.java index ff2d24c..67f064c 100644 --- a/workflow-service/src/main/java/com/flowforge/workflowservice/domain/execution/TaskExecution.java +++ b/workflow-service/src/main/java/com/flowforge/workflowservice/domain/execution/TaskExecution.java @@ -1,7 +1,11 @@ package com.flowforge.workflowservice.domain.execution; +import com.fasterxml.jackson.databind.JsonNode; import com.flowforge.workflowservice.domain.node.WorkflowNode; import jakarta.persistence.*; import lombok.*; +import org.hibernate.annotations.JdbcTypeCode; +import org.hibernate.type.SqlTypes; + import java.time.Instant; import java.util.UUID; @@ -45,4 +49,9 @@ public class TaskExecution { private Integer retryCount = 0; private String errorMessage; + + @JdbcTypeCode(SqlTypes.JSON) + @Column(columnDefinition = "jsonb") + private JsonNode output; + } \ No newline at end of file diff --git a/workflow-service/src/main/java/com/flowforge/workflowservice/domain/node/NodeType.java b/workflow-service/src/main/java/com/flowforge/workflowservice/domain/node/NodeType.java index 7964240..0d55bd2 100644 --- a/workflow-service/src/main/java/com/flowforge/workflowservice/domain/node/NodeType.java +++ b/workflow-service/src/main/java/com/flowforge/workflowservice/domain/node/NodeType.java @@ -7,6 +7,7 @@ public enum NodeType { TASK, EMAIL, HTTP_REQUEST, + SLACK, CONDITION, APPROVAL, diff --git a/workflow-service/src/main/java/com/flowforge/workflowservice/infrastructure/persistence/WorkflowEdgeRepository.java b/workflow-service/src/main/java/com/flowforge/workflowservice/infrastructure/persistence/WorkflowEdgeRepository.java index 7baa8cb..275a90c 100644 --- a/workflow-service/src/main/java/com/flowforge/workflowservice/infrastructure/persistence/WorkflowEdgeRepository.java +++ b/workflow-service/src/main/java/com/flowforge/workflowservice/infrastructure/persistence/WorkflowEdgeRepository.java @@ -1,10 +1,12 @@ package com.flowforge.workflowservice.infrastructure.persistence; +import com.flowforge.workflowservice.domain.edge.BranchType; import com.flowforge.workflowservice.domain.edge.WorkflowEdge; import com.flowforge.workflowservice.domain.node.WorkflowNode; import org.springframework.data.jpa.repository.JpaRepository; import org.springframework.stereotype.Repository; import java.util.List; +import java.util.Optional; import java.util.UUID; @Repository @@ -12,4 +14,8 @@ public interface WorkflowEdgeRepository extends JpaRepository findByWorkflow_Id(UUID workflowId); void deleteByWorkflow_Id(UUID workflowId); List findBySourceNode(WorkflowNode sourceNode); + Optional findBySourceNodeAndBranchType( + WorkflowNode sourceNode, + BranchType branchType + ); } diff --git a/workflow-service/src/main/java/com/flowforge/workflowservice/presentation/dto/request/ConditionConfig.java b/workflow-service/src/main/java/com/flowforge/workflowservice/presentation/dto/request/ConditionConfig.java new file mode 100644 index 0000000..3ecc841 --- /dev/null +++ b/workflow-service/src/main/java/com/flowforge/workflowservice/presentation/dto/request/ConditionConfig.java @@ -0,0 +1,11 @@ +package com.flowforge.workflowservice.presentation.dto.request; + +import com.fasterxml.jackson.databind.JsonNode; +import com.flowforge.workflowservice.application.execution.condtion.ConditionOperator; + +public record ConditionConfig( + String field, + ConditionOperator operator, + JsonNode value + +) {} \ No newline at end of file diff --git a/workflow-service/src/main/resources/application.properties b/workflow-service/src/main/resources/application.properties index 20dd496..7db992e 100644 --- a/workflow-service/src/main/resources/application.properties +++ b/workflow-service/src/main/resources/application.properties @@ -36,9 +36,6 @@ spring.kafka.consumer.enable-auto-commit=false spring.kafka.consumer.key-deserializer=org.apache.kafka.common.serialization.StringDeserializer spring.kafka.consumer.value-deserializer=org.springframework.kafka.support.serializer.JsonDeserializer spring.kafka.consumer.properties.spring.json.trusted.packages=* -spring.kafka.consumer.properties.spring.json.use.type.headers=false -spring.kafka.consumer.properties.spring.json.value.default.type=com.flowforge.workflowservice.application.execution.event.TaskCompletedEvent -spring.kafka.consumer.properties.spring.deserializer.value.delegate.class=org.springframework.kafka.support.serializer.JsonDeserializer # Listener spring.kafka.listener.ack-mode=record