Remove sink blocks and harden execution validation

This commit is contained in:
Lucio Lelii 2026-03-09 15:33:33 +01:00
parent 26e5077a98
commit cb39b2ec63
12 changed files with 90 additions and 55 deletions

View File

@ -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<T extends BlockType> {
final String id = UUID.randomUUID().toString();
Position position;
boolean sink;
String name;
List<IODescriptor> inputs;
@ -64,7 +64,4 @@ public class Block<T extends BlockType> {
this.type = type;
}
public void asSink() {
this.sink = true;
}
}

View File

@ -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<Human
private String actionDescription;
@NotNull
@Valid
@JsonProperty(required = true)
private LLMDescriptor simulateWith;

View File

@ -4,6 +4,7 @@ import com.fasterxml.jackson.annotation.JsonProperty;
import it.cnr.isti.workflow.manager.blocks.types.LLMBlockType;
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;
@ -20,6 +21,7 @@ public class LLMBlockConfiguration extends BlockConfiguration<LLMBlockType> {
@NotNull
@Valid
@JsonProperty(required = true)
LLMDescriptor llmDescriptor;

View File

@ -70,14 +70,20 @@ public class ExecutionContext implements ExecutionListener {
@Override
public void completed(String id, Map<String, Object> 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);
}
}

View File

@ -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;
}

View File

@ -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) {}

View File

@ -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<GenerateResponse> result = webClient.post()
Mono<String> 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<String> getRegisteredModels() {

View File

@ -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<LLMBlockType> draftBlock = blocksController.create(LLMBlockConfiguration.empty());
draftBlock.asSink();
Flow flow = Flow.builder()
.name("Draft Flow")
.description("Flow saved before configuration is complete")

View File

@ -326,8 +326,6 @@ public class FlowControllerTest {
@Test
public void createDraftFlowReturnsDraftStatus() {
Block<LLMBlockType> 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<LLMBlockType> 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<FlowView> createResponse = flowController.createFlow(
request,
new LoginEntity("testuser", "testpassword"));
assertTrue(createResponse.getStatusCode().is2xxSuccessful());
assertNotNull(createResponse.getBody());
assertEquals(FlowViewStatus.DRAFT, createResponse.getBody().status());
}
}

View File

@ -52,8 +52,6 @@ public class FlowImportComponentTest {
.build())
.build()))
.build();
flowData.getBlocks().get(0).asSink();
ImportedFlow importedFlow = new ImportedFlow(
"seed-flow-id",
"Imported Flow",

View File

@ -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())

View File

@ -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"
}
]
}
}