From cb39b2ec6379b803f586747885dbd58979f7c890 Mon Sep 17 00:00:00 2001 From: Lucio Lelii Date: Mon, 9 Mar 2026 15:33:33 +0100 Subject: [PATCH] Remove sink blocks and harden execution validation --- .../isti/workflow/manager/blocks/Block.java | 7 ++-- .../HumanInteractiveBlockConfiguration.java | 2 ++ .../configurations/LLMBlockConfiguration.java | 2 ++ .../manager/executions/ExecutionContext.java | 22 +++++++----- .../manager/executions/ExecutionObject.java | 7 ---- .../workflow/manager/llms/LLMDescriptor.java | 6 ++-- .../ollama/InternalOllamaLLMProvider.java | 28 +++++++++++++-- .../controllers/ExecutionControllerTest.java | 4 --- .../controllers/FlowControllerTest.java | 35 +++++++++++++++++-- .../flows/FlowImportComponentTest.java | 2 -- .../manager/flows/FlowTestCreator.java | 6 ---- src/test/resources/interactionFlow.json | 24 +++++-------- 12 files changed, 90 insertions(+), 55 deletions(-) diff --git a/src/main/java/it/cnr/isti/workflow/manager/blocks/Block.java b/src/main/java/it/cnr/isti/workflow/manager/blocks/Block.java index a33ce30..528bb70 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/blocks/Block.java +++ b/src/main/java/it/cnr/isti/workflow/manager/blocks/Block.java @@ -3,6 +3,7 @@ package it.cnr.isti.workflow.manager.blocks; import java.util.List; import java.util.UUID; +import com.fasterxml.jackson.annotation.JsonIgnoreProperties; import com.fasterxml.jackson.annotation.JsonIgnore; import com.fasterxml.jackson.annotation.JsonProperty; @@ -21,14 +22,13 @@ import lombok.ToString; @NoArgsConstructor(access = lombok.AccessLevel.PROTECTED) @Getter @ToString +@JsonIgnoreProperties(ignoreUnknown = true) public class Block { final String id = UUID.randomUUID().toString(); Position position; - boolean sink; - String name; List inputs; @@ -64,7 +64,4 @@ public class Block { this.type = type; } - public void asSink() { - this.sink = true; - } } diff --git a/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/HumanInteractiveBlockConfiguration.java b/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/HumanInteractiveBlockConfiguration.java index 46971c0..0a920ad 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/HumanInteractiveBlockConfiguration.java +++ b/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/HumanInteractiveBlockConfiguration.java @@ -4,6 +4,7 @@ import com.fasterxml.jackson.annotation.JsonProperty; import it.cnr.isti.workflow.manager.blocks.types.HumanInteractionBlockType; import it.cnr.isti.workflow.manager.llms.LLMDescriptor; +import jakarta.validation.Valid; import jakarta.validation.constraints.NotBlank; import jakarta.validation.constraints.NotNull; import lombok.Builder; @@ -24,6 +25,7 @@ public class HumanInteractiveBlockConfiguration extends BlockConfiguration { @NotNull + @Valid @JsonProperty(required = true) LLMDescriptor llmDescriptor; diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionContext.java b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionContext.java index 142f3bd..4debaea 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionContext.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionContext.java @@ -70,14 +70,20 @@ public class ExecutionContext implements ExecutionListener { @Override public void completed(String id, Map result) { logger.info("Step " + id + " completed with result: " + result); - if (this.steps.get(id).getBlock().isSink()) { - logger.info("All inputs for step " + id + " that is a sink have been satisfied."); - result.forEach((k, v) -> this.addResult(id, k, v)); - boolean allCompleted = this.steps.values().stream().filter(s -> s.getBlock().isSink()) - .allMatch(s -> s.getStatus() == StepStatus.COMPLETED); - if (allCompleted) { - this.setStatus(ExecutionStatus.SUCCESS); - } + Step completedStep = this.steps.get(id); + completedStep.getOutputs().stream() + .filter(output -> !output.isConnected()) + .forEach(output -> { + String outputName = output.getDescriptor().getName(); + if (result.containsKey(outputName)) { + this.addResult(id, outputName, result.get(outputName)); + } + }); + + boolean allCompleted = this.steps.values().stream() + .allMatch(step -> step.getStatus() == StepStatus.COMPLETED); + if (allCompleted) { + this.setStatus(ExecutionStatus.SUCCESS); } } diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionObject.java b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionObject.java index 9aa5de1..081097e 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionObject.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionObject.java @@ -61,13 +61,6 @@ public class ExecutionObject { Input input = targetStep.getInputs().stream().filter(i -> i.getDescriptor().getName().equals(connection.getTargetName())).findFirst().orElseThrow(); output.registerSource(input); } - for (Step step : steps) { - for (Output output : step.getOutputs()) { - if (!output.isConnected() && !step.getBlock().isSink()) - throw new IllegalStateException("Output " + output.getDescriptor().getName() + " of step " + step.getId() - + " is not connected to any input and the block is not a sink"); - } - } return steps; } diff --git a/src/main/java/it/cnr/isti/workflow/manager/llms/LLMDescriptor.java b/src/main/java/it/cnr/isti/workflow/manager/llms/LLMDescriptor.java index a92b9d1..467cf68 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/llms/LLMDescriptor.java +++ b/src/main/java/it/cnr/isti/workflow/manager/llms/LLMDescriptor.java @@ -1,10 +1,10 @@ package it.cnr.isti.workflow.manager.llms; import it.cnr.isti.workflow.manager.blocks.configurations.FieldRetriever; -import jakarta.validation.constraints.NotNull; +import jakarta.validation.constraints.NotBlank; import lombok.Builder; @Builder public record LLMDescriptor( - @NotNull @FieldRetriever(name = "LLM", url = "/retriever/LLM/providers") String provider, - @NotNull @FieldRetriever(name = "LLM", url = "/retriever/LLM/models", dependsOn = {"provider"}) String model) {} + @NotBlank @FieldRetriever(name = "LLM", url = "/retriever/LLM/providers") String provider, + @NotBlank @FieldRetriever(name = "LLM", url = "/retriever/LLM/models", dependsOn = {"provider"}) String model) {} diff --git a/src/main/java/it/cnr/isti/workflow/manager/llms/providers/ollama/InternalOllamaLLMProvider.java b/src/main/java/it/cnr/isti/workflow/manager/llms/providers/ollama/InternalOllamaLLMProvider.java index f8973e4..7eb1c67 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/llms/providers/ollama/InternalOllamaLLMProvider.java +++ b/src/main/java/it/cnr/isti/workflow/manager/llms/providers/ollama/InternalOllamaLLMProvider.java @@ -63,7 +63,7 @@ public class InternalOllamaLLMProvider implements LLMProvider { // Implement the logic to call the Ollama API and return the response WebClient webClient = webClientBuilder.baseUrl(this.ollamaURL).build(); - Mono result = webClient.post() + Mono result = webClient.post() .uri(uriBuilder -> uriBuilder.pathSegment("generate") .build()) .header("Authorization", "Bearer " + ollamaKey) @@ -76,10 +76,32 @@ public class InternalOllamaLLMProvider implements LLMProvider { .defaultIfEmpty("Error: server without body") .flatMap(body -> Mono.error(new RuntimeException( "Error 5xx: " + body)))) - .bodyToMono(GenerateResponse.class) // deserialize JSON in oggetto Java + .bodyToMono(String.class) .timeout(Duration.ofMinutes(2)); - return result.block().getResponse(); // Placeholder for actual implementation + return parseGenerateResponse(result.block(), mapper); + } + + private String parseGenerateResponse(String responseBody, ObjectMapper mapper) { + if (responseBody == null || responseBody.isBlank()) { + throw new RuntimeException("Empty response body from Ollama generate endpoint"); + } + + String trimmedResponse = responseBody.trim(); + if (trimmedResponse.startsWith("{")) { + try { + GenerateResponse response = mapper.readValue(trimmedResponse, GenerateResponse.class); + if (response.getResponse() == null) { + throw new RuntimeException("Missing 'response' field in Ollama JSON response"); + } + return response.getResponse(); + } catch (Exception e) { + throw new RuntimeException("Unable to parse Ollama JSON response", e); + } + } + + log.debug("Ollama generate endpoint returned text/plain response"); + return trimmedResponse; } public List getRegisteredModels() { diff --git a/src/test/java/it/cnr/isti/workflow/manager/controllers/ExecutionControllerTest.java b/src/test/java/it/cnr/isti/workflow/manager/controllers/ExecutionControllerTest.java index aa83d49..4bb5663 100644 --- a/src/test/java/it/cnr/isti/workflow/manager/controllers/ExecutionControllerTest.java +++ b/src/test/java/it/cnr/isti/workflow/manager/controllers/ExecutionControllerTest.java @@ -59,8 +59,6 @@ public class ExecutionControllerTest { .name("second") .llmDescriptor(llmDescriptor) .build()); - second.asSink(); - Connection connection = Connection.builder() .sourceId(first.getId()) .sourceName(first.getOutputs().get(0).getName()) @@ -123,8 +121,6 @@ public class ExecutionControllerTest { @Test public void createExecutionRejectsDraftFlowWithMissingRequiredConfiguration() { Block draftBlock = blocksController.create(LLMBlockConfiguration.empty()); - draftBlock.asSink(); - Flow flow = Flow.builder() .name("Draft Flow") .description("Flow saved before configuration is complete") diff --git a/src/test/java/it/cnr/isti/workflow/manager/controllers/FlowControllerTest.java b/src/test/java/it/cnr/isti/workflow/manager/controllers/FlowControllerTest.java index 262a309..98605e6 100644 --- a/src/test/java/it/cnr/isti/workflow/manager/controllers/FlowControllerTest.java +++ b/src/test/java/it/cnr/isti/workflow/manager/controllers/FlowControllerTest.java @@ -326,8 +326,6 @@ public class FlowControllerTest { @Test public void createDraftFlowReturnsDraftStatus() { Block draftBlock = blocksController.create(LLMBlockConfiguration.empty()); - draftBlock.asSink(); - Flow draftFlow = Flow.builder() .name("Draft Flow") .description("Incomplete but structurally valid flow") @@ -350,4 +348,37 @@ public class FlowControllerTest { assertNotNull(createResponse.getBody()); assertEquals(FlowViewStatus.DRAFT, createResponse.getBody().status()); } + + @Test + public void createFlowWithBlankModelReturnsDraftStatus() { + Block draftBlock = blocksController.create(LLMBlockConfiguration.builder() + .name("translation") + .prompt("translate ${{phrase}} in ${{language}}, return only the translation") + .llmDescriptor(LLMDescriptor.builder() + .provider("Gemini") + .model("") + .build()) + .build()); + Flow draftFlow = Flow.builder() + .name("translation flow") + .description("") + .block(draftBlock) + .build(); + + FlowCreateRequest request = new FlowCreateRequest( + draftFlow.getName(), + draftFlow.getDescription(), + FlowData.builder() + .blocks(draftFlow.getBlocks()) + .connections(draftFlow.getConnections()) + .build()); + + ResponseEntity createResponse = flowController.createFlow( + request, + new LoginEntity("testuser", "testpassword")); + + assertTrue(createResponse.getStatusCode().is2xxSuccessful()); + assertNotNull(createResponse.getBody()); + assertEquals(FlowViewStatus.DRAFT, createResponse.getBody().status()); + } } diff --git a/src/test/java/it/cnr/isti/workflow/manager/flows/FlowImportComponentTest.java b/src/test/java/it/cnr/isti/workflow/manager/flows/FlowImportComponentTest.java index 781c1d9..6291499 100644 --- a/src/test/java/it/cnr/isti/workflow/manager/flows/FlowImportComponentTest.java +++ b/src/test/java/it/cnr/isti/workflow/manager/flows/FlowImportComponentTest.java @@ -52,8 +52,6 @@ public class FlowImportComponentTest { .build()) .build())) .build(); - flowData.getBlocks().get(0).asSink(); - ImportedFlow importedFlow = new ImportedFlow( "seed-flow-id", "Imported Flow", diff --git a/src/test/java/it/cnr/isti/workflow/manager/flows/FlowTestCreator.java b/src/test/java/it/cnr/isti/workflow/manager/flows/FlowTestCreator.java index 2728878..d52fa6c 100644 --- a/src/test/java/it/cnr/isti/workflow/manager/flows/FlowTestCreator.java +++ b/src/test/java/it/cnr/isti/workflow/manager/flows/FlowTestCreator.java @@ -41,8 +41,6 @@ public class FlowTestCreator { .name("second") .llmDescriptor(llmBrick) .build()); - block2.asSink(); - Connection connection = Connection.builder() .sourceId(block1.getId()) .sourceName(block1.getOutputs().get(0).getName()) @@ -69,8 +67,6 @@ public class FlowTestCreator { .name("second") .llmDescriptor(llmBrick) .build()); - block2.asSink(); - Connection connection = Connection.builder() .sourceId(block1.getId()) .sourceName(block1.getOutputs().get(0).getName()) @@ -97,8 +93,6 @@ public class FlowTestCreator { .simulateWith(llmBrick) .build()); - block2.asSink(); - Connection connection = Connection.builder() .sourceId(block1.getId()) .sourceName(block1.getOutputs().get(0).getName()) diff --git a/src/test/resources/interactionFlow.json b/src/test/resources/interactionFlow.json index 96e95e7..1c5144d 100644 --- a/src/test/resources/interactionFlow.json +++ b/src/test/resources/interactionFlow.json @@ -4,7 +4,6 @@ "blocks": [ { "id": "d1f8bca5-5976-489d-8cc9-49b062c5e2bf", - "sink": false, "name": "first", "inputs": [ { @@ -23,19 +22,16 @@ "specificConfiguration": { "type": "LLMBlockConfiguration", "name": "first", - "prompt": "Hello, ${{name}}! How are you ${{name}}?", - "brick": { - "id": "testBrick:llmTestProvider_testModel", - "brickManager": "testBrick", - "model": "testModel", - "executionTimeConfiguration": {} - } + "llmDescriptor": { + "provider": "testProvider", + "model": "testModel" + }, + "prompt": "Hello, ${{name}}! How are you ${{name}}?" }, "typeName": "LLMBlock" }, { "id": "0f7852c5-c11a-4dda-9b2d-9cb76e58e55a", - "sink": true, "name": "interactive", "inputs": [ { @@ -55,11 +51,9 @@ "type": "HumanInteractiveBlockConfiguration", "name": "interactive", "actionDescription": "Answer the question", - "llmSimulatorBrick": { - "id": "testBrick:humanInteractionTestProvider_testModel", - "brickManager": "testBrick", - "model": "testModel", - "executionTimeConfiguration": {} + "simulateWith": { + "provider": "testProvider", + "model": "testModel" } }, "typeName": "HumanInteractionBlock" @@ -74,4 +68,4 @@ "targetName": "input" } ] -} \ No newline at end of file +}