changed Flow structure

This commit is contained in:
Lucio Lelii 2026-03-06 12:06:07 +01:00
parent d0e4823a70
commit d7f894acb4
16 changed files with 288 additions and 45 deletions

View File

@ -101,7 +101,7 @@ public class ExecutionsController {
FlowEntity flow = flowRepository.findById(flowId)
.orElseThrow(() -> new IllegalArgumentException("Flow with id " + flowId + " not found"));
try{
ExecutionObject eo = executionService.createExecution(flow.getFlow());
ExecutionObject eo = executionService.createExecution(flow.getName(), flow.getFlow());
return eo;
} catch (Throwable e) {
logger.error("Error creating execution for flow {}", flowId, e);

View File

@ -5,17 +5,18 @@ import java.util.List;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.http.ResponseEntity;
import org.springframework.security.core.annotation.AuthenticationPrincipal;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
import org.eclipse.microprofile.openapi.annotations.Operation;
import jakarta.validation.Valid;
import it.cnr.isti.workflow.manager.auth.repo.LoginEntity;
import it.cnr.isti.workflow.manager.flows.FlowService;
import it.cnr.isti.workflow.manager.flows.model.Flow;
import it.cnr.isti.workflow.manager.flows.repo.FlowEntity;
import org.springframework.web.bind.annotation.GetMapping;
import it.cnr.isti.workflow.manager.flows.model.FlowCreateRequest;
import it.cnr.isti.workflow.manager.flows.model.FlowView;
@RestController
@ -27,7 +28,7 @@ public class FlowController {
@GetMapping("")
@Operation(summary = "Get all flows", description = "Returns the list of all available workflow definitions.")
public List<Flow> getAllFlows() {
public List<FlowView> getAllFlows() {
return flowService.getAllFlows();
}
@ -35,9 +36,10 @@ public class FlowController {
@PostMapping
@Operation(summary = "Create flow", description = "Creates a new flow owned by the authenticated user.")
public ResponseEntity<FlowEntity> createFlow(@RequestBody Flow flow, @AuthenticationPrincipal
public ResponseEntity<FlowView> createFlow(@RequestBody @Valid FlowCreateRequest flow,
@AuthenticationPrincipal
LoginEntity userDetails) {
FlowEntity createdFlow = flowService.createFlow(userDetails.getUsername(), flow);
FlowView createdFlow = flowService.createFlow(userDetails.getUsername(), flow);
return ResponseEntity.ok(createdFlow);
}

View File

@ -15,7 +15,7 @@ import it.cnr.isti.workflow.manager.executions.steps.Input;
import it.cnr.isti.workflow.manager.executions.steps.Output;
import it.cnr.isti.workflow.manager.executions.steps.Step;
import it.cnr.isti.workflow.manager.flows.model.Connection;
import it.cnr.isti.workflow.manager.flows.model.Flow;
import it.cnr.isti.workflow.manager.flows.model.FlowData;
import lombok.Builder;
import lombok.Getter;
import lombok.NoArgsConstructor;
@ -36,7 +36,7 @@ public class ExecutionObject {
ExecutorService executorService;
@Builder
public ExecutionObject(String executionName, Flow flow) {
public ExecutionObject(String executionName, FlowData flow) {
this.name = executionName;
List<Step<?>> steps = getStepsFromFlow(flow);
@ -47,11 +47,11 @@ public class ExecutionObject {
}
List<Step<?>> getStepsFromFlow(Flow flow) {
List<Step<?>> getStepsFromFlow(FlowData flow) {
List<Step<?>> steps = new ArrayList<>();
flow.getBlocks().stream().filter(b -> BlockExecutors.get(b.getType()) == null).findAny().ifPresent(b -> {
throw new IllegalStateException("No executor found for block type " + b.getType().getName()
+ ". Cannot create execution for flow " + flow.getName());
+ ". Cannot create execution");
});
flow.getBlocks().forEach(block -> steps.add(new Step<>(block)));
for (Connection connection: flow.getConnections()){
@ -96,4 +96,4 @@ public class ExecutionObject {
+ " is not in READY status (CURRENT STATUS is " + this.getContext().getStatus() + ")");
}
}
}

View File

@ -5,6 +5,7 @@ import java.util.List;
import java.util.Map;
import org.springframework.stereotype.Service;
import it.cnr.isti.workflow.manager.flows.model.Flow;
import it.cnr.isti.workflow.manager.flows.model.FlowData;
@Service
public class ExecutionsService {
@ -12,15 +13,23 @@ public class ExecutionsService {
private static Map<String, ExecutionObject> executions = new HashMap<>();
public ExecutionObject createExecution(Flow flow) {
public ExecutionObject createExecution(String executionName, FlowData flow) {
ExecutionObject execObject = ExecutionObject.builder()
.executionName(flow.getName())
.executionName(executionName)
.flow(flow)
.build();
executions.put(execObject.getId(), execObject);
return execObject;
}
public ExecutionObject createExecution(Flow flow) {
FlowData flowData = FlowData.builder()
.blocks(flow.getBlocks())
.connections(flow.getConnections())
.build();
return createExecution(flow.getName(), flowData);
}
public ExecutionObject getExecution(String id) {
ExecutionObject toReturn = executions.get(id);
if (toReturn == null)
@ -59,4 +68,4 @@ public class ExecutionsService {
}
}
}

View File

@ -2,36 +2,36 @@ package it.cnr.isti.workflow.manager.flows;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;
import it.cnr.isti.workflow.manager.flows.model.Flow;
import it.cnr.isti.workflow.manager.flows.model.FlowData;
import jakarta.persistence.AttributeConverter;
import jakarta.persistence.Converter;
@Converter(autoApply = false)
public class FlowConverter implements AttributeConverter<Flow, String> {
public class FlowConverter implements AttributeConverter<FlowData, String> {
private static final ObjectMapper objectMapper = new ObjectMapper();
@Override
public String convertToDatabaseColumn(Flow flow) {
if (flow == null) {
public String convertToDatabaseColumn(FlowData flowData) {
if (flowData == null) {
return null;
}
try {
return objectMapper.writeValueAsString(flow);
return objectMapper.writeValueAsString(flowData);
} catch (JsonProcessingException e) {
throw new IllegalArgumentException("Errore nella serializzazione di Flow in JSON", e);
throw new IllegalArgumentException("Errore nella serializzazione di FlowData in JSON", e);
}
}
@Override
public Flow convertToEntityAttribute(String dbData) {
public FlowData convertToEntityAttribute(String dbData) {
if (dbData == null || dbData.isBlank()) {
return null;
}
try {
return objectMapper.readValue(dbData, Flow.class);
return objectMapper.readValue(dbData, FlowData.class);
} catch (Exception e) {
throw new IllegalArgumentException("Errore nella deserializzazione di JSON in Flow", e);
throw new IllegalArgumentException("Errore nella deserializzazione di JSON in FlowData", e);
}
}
}

View File

@ -0,0 +1,38 @@
package it.cnr.isti.workflow.manager.flows;
import java.time.LocalDateTime;
import it.cnr.isti.workflow.manager.flows.model.FlowCreateRequest;
import it.cnr.isti.workflow.manager.flows.model.FlowView;
import it.cnr.isti.workflow.manager.flows.repo.FlowEntity;
public final class FlowMapper {
private FlowMapper() {
}
public static FlowEntity toNewEntity(String owner, FlowCreateRequest request) {
LocalDateTime now = LocalDateTime.now();
FlowEntity entity = new FlowEntity();
entity.setName(request.name());
entity.setDescription(request.description());
entity.setOwner(owner);
entity.setCreatedAt(now);
entity.setLastUpdateAt(now);
entity.setFlow(request.flow());
return entity;
}
public static FlowView toView(FlowEntity entity) {
return new FlowView(
entity.getId(),
entity.getName(),
entity.getDescription(),
entity.getCreatedAt(),
entity.getLastUpdateAt(),
entity.getOwner(),
entity.isPublished(),
entity.isFinalized(),
entity.getFlow());
}
}

View File

@ -1,12 +1,12 @@
package it.cnr.isti.workflow.manager.flows;
import java.time.LocalDateTime;
import java.util.List;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import it.cnr.isti.workflow.manager.flows.model.Flow;
import it.cnr.isti.workflow.manager.flows.model.FlowCreateRequest;
import it.cnr.isti.workflow.manager.flows.model.FlowView;
import it.cnr.isti.workflow.manager.flows.repo.FlowEntity;
import it.cnr.isti.workflow.manager.flows.repo.FlowRepository;
@ -16,18 +16,13 @@ public class FlowService {
@Autowired
FlowRepository flowRepository;
public FlowEntity createFlow(String owner, Flow flow) {
FlowEntity flowEntity = new FlowEntity();
flowEntity.setName(flow.getName());
LocalDateTime now = LocalDateTime.now();
flowEntity.setCreatedAt(now);
flowEntity.setLastUpdateAt(now);
flowEntity.setFlow(flow);
flowEntity.setOwner(owner);
return flowRepository.save(flowEntity);
public FlowView createFlow(String owner, FlowCreateRequest request) {
FlowEntity flowEntity = FlowMapper.toNewEntity(owner, request);
FlowEntity savedEntity = flowRepository.save(flowEntity);
return FlowMapper.toView(savedEntity);
}
public List<Flow> getAllFlows() {
return flowRepository.findAll().stream().map(FlowEntity::getFlow).toList();
public List<FlowView> getAllFlows() {
return flowRepository.findAll().stream().map(FlowMapper::toView).toList();
}
}

View File

@ -28,5 +28,4 @@ public class Flow {
@Singular
List<Connection> connections;
}

View File

@ -0,0 +1,10 @@
package it.cnr.isti.workflow.manager.flows.model;
import jakarta.validation.constraints.NotBlank;
import jakarta.validation.constraints.NotNull;
public record FlowCreateRequest(
@NotBlank String name,
String description,
@NotNull FlowData flow) {
}

View File

@ -0,0 +1,26 @@
package it.cnr.isti.workflow.manager.flows.model;
import java.util.List;
import it.cnr.isti.workflow.manager.blocks.Block;
import jakarta.validation.constraints.NotBlank;
import lombok.AllArgsConstructor;
import lombok.Builder;
import lombok.Data;
import lombok.NoArgsConstructor;
import lombok.NonNull;
import lombok.Singular;
@Data
@Builder
@NoArgsConstructor(access = lombok.AccessLevel.PROTECTED)
@AllArgsConstructor
public class FlowData {
@Singular
List<Block<?>> blocks;
@Singular
List<Connection> connections;
}

View File

@ -0,0 +1,15 @@
package it.cnr.isti.workflow.manager.flows.model;
import java.time.LocalDateTime;
public record FlowView(
String id,
String name,
String description,
LocalDateTime createdAt,
LocalDateTime lastUpdateAt,
String owner,
boolean published,
boolean finalized,
FlowData flow) {
}

View File

@ -3,7 +3,7 @@ package it.cnr.isti.workflow.manager.flows.repo;
import java.time.LocalDateTime;
import it.cnr.isti.workflow.manager.flows.FlowConverter;
import it.cnr.isti.workflow.manager.flows.model.Flow;
import it.cnr.isti.workflow.manager.flows.model.FlowData;
import jakarta.persistence.Column;
import jakarta.persistence.Convert;
import jakarta.persistence.Entity;
@ -54,5 +54,5 @@ public class FlowEntity {
@Lob
@Column(name = "flow_data", columnDefinition = "TEXT")
@Convert(converter = FlowConverter.class)
private Flow flow;
private FlowData flow;
}

View File

@ -19,7 +19,9 @@ import it.cnr.isti.workflow.manager.executions.ExecutionObject;
import it.cnr.isti.workflow.manager.executions.ExecutionStatus;
import it.cnr.isti.workflow.manager.flows.model.Connection;
import it.cnr.isti.workflow.manager.flows.model.Flow;
import it.cnr.isti.workflow.manager.flows.repo.FlowEntity;
import it.cnr.isti.workflow.manager.flows.model.FlowCreateRequest;
import it.cnr.isti.workflow.manager.flows.model.FlowData;
import it.cnr.isti.workflow.manager.flows.model.FlowView;
import it.cnr.isti.workflow.manager.llms.LLMDescriptor;
@SpringBootTest
@ -72,8 +74,12 @@ public class ExecutionControllerTest {
.connection(connection)
.build();
ResponseEntity<FlowEntity> createdFlow = flowController.createFlow(flow, new LoginEntity("testuser", "testpassword"));
ExecutionObject executionObject = executionsController.create(createdFlow.getBody().getId());
FlowCreateRequest request = new FlowCreateRequest(
flow.getName(),
flow.getDescription(),
FlowData.builder().blocks(flow.getBlocks()).connections(flow.getConnections()).build());
ResponseEntity<FlowView> createdFlow = flowController.createFlow(request, new LoginEntity("testuser", "testpassword"));
ExecutionObject executionObject = executionsController.create(createdFlow.getBody().id());
assert executionObject != null;
assert executionObject.getId() != null;

View File

@ -0,0 +1,84 @@
package it.cnr.isti.workflow.manager.controllers;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
import java.util.List;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.http.ResponseEntity;
import org.springframework.test.context.TestPropertySource;
import com.fasterxml.jackson.core.JsonProcessingException;
import it.cnr.isti.workflow.manager.app.ObjectMapperHolder;
import it.cnr.isti.workflow.manager.auth.repo.LoginEntity;
import it.cnr.isti.workflow.manager.flows.FlowTestCreator;
import it.cnr.isti.workflow.manager.flows.model.Flow;
import it.cnr.isti.workflow.manager.flows.model.FlowCreateRequest;
import it.cnr.isti.workflow.manager.flows.model.FlowData;
import it.cnr.isti.workflow.manager.flows.model.FlowView;
import it.cnr.isti.workflow.manager.llms.LLMDescriptor;
@SpringBootTest
@TestPropertySource(locations = "classpath:test.properties")
public class FlowControllerTest {
@Autowired
private FlowController flowController;
@Autowired
private FlowTestCreator flowTestCreator;
@Test
public void createAndGetFlow() {
LLMDescriptor llmDescriptor = LLMDescriptor.builder()
.provider("testProvider")
.model("testModel")
.build();
Flow flow = flowTestCreator.createFlowWithConnection(llmDescriptor);
FlowCreateRequest request = new FlowCreateRequest(
flow.getName(),
flow.getDescription(),
FlowData.builder()
.blocks(flow.getBlocks())
.connections(flow.getConnections())
.build());
ResponseEntity<FlowView> createResponse = flowController.createFlow(
request,
new LoginEntity("testuser", "testpassword"));
assertTrue(createResponse.getStatusCode().is2xxSuccessful());
assertNotNull(createResponse.getBody());
FlowView created = createResponse.getBody();
assertNotNull(created.id());
List<FlowView> flows = flowController.getAllFlows();
assertNotNull(flows);
FlowView retrieved = flows.stream()
.filter(f -> created.id().equals(f.id()))
.findFirst()
.orElse(null);
try {
System.out.println(ObjectMapperHolder.mapper.writerWithDefaultPrettyPrinter().writeValueAsString(retrieved));
} catch (JsonProcessingException e) {
e.printStackTrace();
}
assertNotNull(retrieved);
assertEquals(created.id(), retrieved.id());
assertEquals("testuser", retrieved.owner());
assertEquals(flow.getName(), retrieved.name());
assertEquals(flow.getDescription(), retrieved.description());
assertNotNull(retrieved.flow());
assertEquals(flow.getBlocks().size(), retrieved.flow().getBlocks().size());
assertEquals(flow.getConnections().size(), retrieved.flow().getConnections().size());
}
}

View File

@ -159,4 +159,54 @@ public class ExecutionWithContainer {
}
}
@Test
public void getInteractiveFlowWithStop() {
Flow flow = flowTestCreator.createFlowWithInteraction(llmBrick);
assertNotNull(flow);
ExecutionObject execObject = executionsService.createExecution(flow);
assertNotNull(execObject);
assertEquals(ExecutionStatus.CREATED, execObject.getContext().getStatus());
try {
logger.info("Execution object in CREATED State: {}", ObjectMapperHolder.mapper.writerWithDefaultPrettyPrinter().writeValueAsString(execObject));
} catch (JsonProcessingException e) {
// TODO Auto-generated catch block
e.printStackTrace();
}
for (Step<?> s : execObject.getContext().getSteps().values()) {
if (s.getInputs().stream().anyMatch(i -> i.getDescriptor().getName().equals("name") && !i.isRegistered())) {
executionsService.prepareInput(execObject.getId(), s.getId(), "name", "marie curie");
}
}
assertEquals(ExecutionStatus.READY, execObject.getContext().getStatus());
try {
logger.info("Execution object in READY State: {}", ObjectMapperHolder.mapper.writerWithDefaultPrettyPrinter().writeValueAsString(execObject));
} catch (JsonProcessingException e) {
// TODO Auto-generated catch block
e.printStackTrace();
}
execObject = executionsService.startExecution(execObject.getId());
while (execObject.getContext().getStatus() == ExecutionStatus.RUNNING) {
try {
Thread.sleep(500);
logger.debug("Execution status: " + execObject.getContext().getStatus());
execObject = executionsService.getExecution(execObject.getId());
} catch (InterruptedException e) {
e.printStackTrace();
}
}
assertEquals(ExecutionStatus.WAITING, execObject.getContext().getStatus());
try {
logger.info("Execution object in WAITING State: {}", ObjectMapperHolder.mapper.writerWithDefaultPrettyPrinter().writeValueAsString(execObject));
} catch (JsonProcessingException e) {
// TODO Auto-generated catch block
e.printStackTrace();
}
}
}

View File

@ -7,6 +7,8 @@ import org.springframework.test.context.TestPropertySource;
import com.fasterxml.jackson.databind.ObjectMapper;
import it.cnr.isti.workflow.manager.flows.model.Flow;
import it.cnr.isti.workflow.manager.flows.model.FlowCreateRequest;
import it.cnr.isti.workflow.manager.flows.model.FlowData;
import it.cnr.isti.workflow.manager.llms.LLMDescriptor;
@SpringBootTest
@ -27,13 +29,13 @@ public class FlowTest {
@Test
public void createEmptyFlow() {
Flow flow = Flow.builder().name("Test Flow").description("This is a test flow").build();
flowService.createFlow("testUser", flow );
flowService.createFlow("testUser", toCreateRequest(flow));
}
@Test
public void createFlow() {
Flow flow = flowTestCreator.createFlowWithConnection(llmBrick);
flowService.createFlow("testUser", flow );
flowService.createFlow("testUser", toCreateRequest(flow));
}
@Test
@ -58,5 +60,12 @@ public class FlowTest {
assert(json.equals(json2));
}
private FlowCreateRequest toCreateRequest(Flow flow) {
FlowData flowData = FlowData.builder()
.blocks(flow.getBlocks())
.connections(flow.getConnections())
.build();
return new FlowCreateRequest(flow.getName(), flow.getDescription(), flowData);
}
}