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