From d3b5fcd9fdd8a0353aeeeee0703b268c8ed4bf9b Mon Sep 17 00:00:00 2001 From: Lucio Lelii Date: Mon, 30 Mar 2026 17:07:13 +0200 Subject: [PATCH] Add global flow inputs and execution gating --- .../blocks/factories/PlaceholderInputs.java | 1 + .../manager/executions/ExecutionContext.java | 14 +++--- .../manager/executions/ExecutionObject.java | 47 ++++++++++++++++++- .../executions/ExecutionTemplateResolver.java | 1 + .../manager/executions/ExecutionsService.java | 3 ++ .../workflow/manager/flows/model/Flow.java | 4 ++ .../manager/flows/model/FlowData.java | 4 ++ .../flows/validation/FlowDataValidator.java | 21 +++++++++ .../manager/executions/ExecutionTest.java | 29 ++++++++++++ 9 files changed, 116 insertions(+), 8 deletions(-) diff --git a/src/main/java/it/cnr/isti/workflow/manager/blocks/factories/PlaceholderInputs.java b/src/main/java/it/cnr/isti/workflow/manager/blocks/factories/PlaceholderInputs.java index 12d7883..fbb5f5c 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/blocks/factories/PlaceholderInputs.java +++ b/src/main/java/it/cnr/isti/workflow/manager/blocks/factories/PlaceholderInputs.java @@ -8,6 +8,7 @@ final class PlaceholderInputs { static boolean isRuntimeExecutionVariable(String placeholder) { return placeholder != null && (placeholder.startsWith("vars.") + || placeholder.startsWith("global.") || placeholder.startsWith("context.")); } } 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 b505222..1579f00 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 @@ -398,13 +398,13 @@ public class ExecutionContext implements ExecutionListener { .endTime(this.endTime) .interactionSimulationEnabled(this.interactionSimulationEnabled) .interactionSimulationDescriptor(this.interactionSimulationDescriptor) - .executionVariables(Map.copyOf(this.executionVariables)) - .executionVariableDescriptors(Map.copyOf(this.executionVariableDescriptors)) - .providedAuthorizations(providedAuthorizations == null ? Map.of() : Map.copyOf(providedAuthorizations)) - .inputs(Map.copyOf(this.inputs)) - .result(Map.copyOf(this.result)) - .partialResult(Map.copyOf(this.partialResult)) - .errors(Map.copyOf(this.errors)) + .executionVariables(new HashMap<>(this.executionVariables)) + .executionVariableDescriptors(new HashMap<>(this.executionVariableDescriptors)) + .providedAuthorizations(providedAuthorizations == null ? Map.of() : new HashMap<>(providedAuthorizations)) + .inputs(new HashMap<>(this.inputs)) + .result(new HashMap<>(this.result)) + .partialResult(new HashMap<>(this.partialResult)) + .errors(new HashMap<>(this.errors)) .warnings(List.copyOf(this.warnings)) .events(List.copyOf(this.events)) .waitingSteps(List.copyOf(this.waitingSteps)) 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 abc6797..87516e4 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 @@ -24,6 +24,7 @@ import it.cnr.isti.workflow.manager.flows.model.FlowNode; import it.cnr.isti.workflow.manager.flows.model.Connection; import it.cnr.isti.workflow.manager.flows.model.Dependency; import it.cnr.isti.workflow.manager.flows.model.FlowData; +import it.cnr.isti.workflow.manager.ios.IODescriptor; import it.cnr.isti.workflow.manager.llms.LLMDescriptor; import lombok.Builder; import lombok.Getter; @@ -52,6 +53,8 @@ public class ExecutionObject { List requiredAuthorizations = List.of(); + List requiredGlobalInputs = List.of(); + Map providedAuthorizations = new HashMap<>(); @JsonIgnore @@ -71,6 +74,7 @@ public class ExecutionObject { this.flow = flow; this.stepConnections = flow.getConnections() == null ? List.of() : List.copyOf(flow.getConnections()); this.stepDependencies = flow.getDependencies() == null ? List.of() : List.copyOf(flow.getDependencies()); + this.requiredGlobalInputs = flow.getGlobalInputs() == null ? List.of() : List.copyOf(flow.getGlobalInputs()); this.requiredAuthorizations = requiredAuthorizations == null ? List.of() : List.copyOf(requiredAuthorizations); List> steps = getStepsFromFlow(flow); @@ -79,6 +83,7 @@ public class ExecutionObject { this.context = new ExecutionContext(steps.stream().collect(Collectors.toMap(Step::getId, Function.identity()))); ensureRuntimeContextVariables(); + registerGlobalInputs(); this.requiredAuthorizations.stream() .map(ExecutionAuthorizationRequirement::key) .distinct() @@ -139,16 +144,22 @@ public class ExecutionObject { protected void setExecutionVariable(String key, Object value) { this.context.setExecutionVariable(key, value); ensureRuntimeContextVariables(); + registerGlobalInputs(); + refreshInitializationStatus(); } protected void setExecutionVariables(Map executionVariables) { this.context.setExecutionVariables(ExecutionRuntimeContextSupport.apply(executionVariables, this.id, this.name, this.owner, this.creationTime)); + registerGlobalInputs(); + refreshInitializationStatus(); } protected void setExecutionVariableDescriptors(Map executionVariableDescriptors) { this.context.setExecutionVariableDescriptors(executionVariableDescriptors); ensureRuntimeContextVariables(); + registerGlobalInputs(); + refreshInitializationStatus(); } protected void registerExecutionVariable(ExecutionVariableDescriptor descriptor) { @@ -158,6 +169,8 @@ public class ExecutionObject { protected Object removeExecutionVariable(String key) { Object removed = this.context.removeExecutionVariable(key); ensureRuntimeContextVariables(); + registerGlobalInputs(); + refreshInitializationStatus(); return removed; } @@ -231,6 +244,16 @@ public class ExecutionObject { .toList(); } + public List getMissingGlobalInputKeys() { + return this.requiredGlobalInputs.stream() + .map(IODescriptor::getName) + .filter(name -> { + ExecutionVariableDescriptor descriptor = this.context.getExecutionVariableDescriptor(name); + return descriptor == null || descriptor.getValue() == null; + }) + .toList(); + } + private void refreshInitializationStatus() { if (!this.context.getStatus().isInitState()) { return; @@ -240,7 +263,7 @@ public class ExecutionObject { || step.getStatus() == it.cnr.isti.workflow.manager.executions.steps.StepStatus.WAITING_FOR_DEPENDENCY || (step.getInputs().stream().allMatch(input -> input.isRegistered() || input.isSet()) && !step.getInputs().stream().anyMatch(input -> !input.isRegistered() && !input.isSet()))); - if (allInputsSatisfied && getMissingAuthorizationKeys().isEmpty()) { + if (allInputsSatisfied && getMissingAuthorizationKeys().isEmpty() && getMissingGlobalInputKeys().isEmpty()) { this.context.setStatus(ExecutionStatus.READY); } else { this.context.setStatus(ExecutionStatus.CREATED); @@ -278,6 +301,7 @@ public class ExecutionObject { this.interactionSimulationDescriptor = snapshot.getInteractionSimulationDescriptor(); } ensureRuntimeContextVariables(); + registerGlobalInputs(); rebuildDependencyStates(); this.providedAuthorizations.forEach(this.context::setAuthorization); this.context.setInteractionSimulationDescriptor(this.interactionSimulationDescriptor); @@ -332,4 +356,25 @@ public class ExecutionObject { .build())); } + private void registerGlobalInputs() { + for (IODescriptor globalInput : this.requiredGlobalInputs) { + ExecutionVariableDescriptor existing = this.context.getExecutionVariableDescriptor(globalInput.getName()); + this.context.registerExecutionVariable(ExecutionVariableDescriptor.builder() + .name(globalInput.getName()) + .value(existing == null ? null : existing.getValue()) + .kind(toExecutionVariableKind(globalInput)) + .description("Global flow input") + .cleanupPolicy(ExecutionVariableCleanupPolicy.NONE) + .build()); + } + } + + private ExecutionVariableKind toExecutionVariableKind(IODescriptor descriptor) { + return switch (descriptor.getType()) { + case TEXT -> ExecutionVariableKind.TEXT; + case FILE, CSV -> ExecutionVariableKind.FILE_PATH; + case ANY -> ExecutionVariableKind.ANY; + }; + } + } diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionTemplateResolver.java b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionTemplateResolver.java index 85968de..e70eadf 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionTemplateResolver.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionTemplateResolver.java @@ -36,6 +36,7 @@ public final class ExecutionTemplateResolver { if (executionVariables != null) { for (Map.Entry entry : executionVariables.entrySet()) { resolved = resolved.replace("${{vars." + entry.getKey() + "}}", formatValue(entry.getValue())); + resolved = resolved.replace("${{global." + entry.getKey() + "}}", formatValue(entry.getValue())); if (entry.getKey() != null && entry.getKey().startsWith(ExecutionRuntimeContextSupport.CONTEXT_PREFIX)) { resolved = resolved.replace("${{" + entry.getKey() + "}}", formatValue(entry.getValue())); } diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionsService.java b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionsService.java index b7b892f..443b670 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionsService.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionsService.java @@ -83,6 +83,9 @@ public class ExecutionsService { if (flow.getDependencies() != null) { flowDataBuilder.dependencies(flow.getDependencies()); } + if (flow.getGlobalInputs() != null) { + flowDataBuilder.globalInputs(flow.getGlobalInputs()); + } FlowData flowData = flowDataBuilder.build(); return createExecution(flow.getName(), flowData, null); } diff --git a/src/main/java/it/cnr/isti/workflow/manager/flows/model/Flow.java b/src/main/java/it/cnr/isti/workflow/manager/flows/model/Flow.java index ae61236..97dbfd3 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/flows/model/Flow.java +++ b/src/main/java/it/cnr/isti/workflow/manager/flows/model/Flow.java @@ -4,6 +4,7 @@ import java.util.List; import it.cnr.isti.workflow.manager.blocks.Block; import it.cnr.isti.workflow.manager.containers.Container; +import it.cnr.isti.workflow.manager.ios.IODescriptor; import jakarta.validation.constraints.NotBlank; import lombok.AllArgsConstructor; import lombok.Builder; @@ -35,4 +36,7 @@ public class Flow { @Singular List dependencies; + + @Singular + List globalInputs; } diff --git a/src/main/java/it/cnr/isti/workflow/manager/flows/model/FlowData.java b/src/main/java/it/cnr/isti/workflow/manager/flows/model/FlowData.java index 8f42aa9..43a1c66 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/flows/model/FlowData.java +++ b/src/main/java/it/cnr/isti/workflow/manager/flows/model/FlowData.java @@ -7,6 +7,7 @@ import com.fasterxml.jackson.annotation.JsonIgnore; import it.cnr.isti.workflow.manager.blocks.Block; import it.cnr.isti.workflow.manager.containers.Container; import it.cnr.isti.workflow.manager.flows.validation.ValidFlowStructure; +import it.cnr.isti.workflow.manager.ios.IODescriptor; import lombok.AllArgsConstructor; import lombok.Builder; import lombok.Data; @@ -33,6 +34,9 @@ public class FlowData { @Singular List dependencies; + @Singular + List globalInputs; + @JsonIgnore public List getNodes() { List> blockList = blocks == null ? List.of() : blocks; diff --git a/src/main/java/it/cnr/isti/workflow/manager/flows/validation/FlowDataValidator.java b/src/main/java/it/cnr/isti/workflow/manager/flows/validation/FlowDataValidator.java index 139fea9..90ad297 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/flows/validation/FlowDataValidator.java +++ b/src/main/java/it/cnr/isti/workflow/manager/flows/validation/FlowDataValidator.java @@ -58,6 +58,9 @@ public class FlowDataValidator implements ConstraintValidator nodes = flowData.getNodes(); List connections = flowData.getConnections() == null ? List.of() : flowData.getConnections(); List dependencies = flowData.getDependencies() == null ? List.of() : flowData.getDependencies(); + List globalInputs = flowData.getGlobalInputs() == null ? List.of() : flowData.getGlobalInputs(); + + validateGlobalInputs(globalInputs); HashMap nodesById = new HashMap<>(); for (FlowNode node : nodes) { @@ -76,6 +79,24 @@ public class FlowDataValidator implements ConstraintValidator globalInputs) { + Set names = new LinkedHashSet<>(); + for (IODescriptor descriptor : globalInputs) { + if (descriptor == null) { + throw validationError(error(ValidationErrorCode.VALIDATION_ERROR, "flow", null, "globalInputs", + "Flow contains a null global input")); + } + if (descriptor.getName() == null || descriptor.getName().isBlank()) { + throw validationError(error(ValidationErrorCode.VALIDATION_ERROR, "flow", null, "globalInputs.name", + "Global input name is required")); + } + if (!names.add(descriptor.getName())) { + throw validationError(error(ValidationErrorCode.VALIDATION_ERROR, "flow", null, "globalInputs", + "Duplicate global input name: " + descriptor.getName())); + } + } + } + private void validateNode(FlowNode node) { if (node == null) { throw validationError(error(ValidationErrorCode.FLOW_CONTAINS_NULL_NODE, "flow", null, "nodes", "Flow contains a null node")); diff --git a/src/test/java/it/cnr/isti/workflow/manager/executions/ExecutionTest.java b/src/test/java/it/cnr/isti/workflow/manager/executions/ExecutionTest.java index 45c10de..a627131 100644 --- a/src/test/java/it/cnr/isti/workflow/manager/executions/ExecutionTest.java +++ b/src/test/java/it/cnr/isti/workflow/manager/executions/ExecutionTest.java @@ -182,6 +182,35 @@ public class ExecutionTest { assertEquals(ExecutionStatus.CREATED, execObject.getContext().getStatus()); } + @Test + public void globalInputsAreRequiredOnceAndDoNotCreateNodeInputs() { + Block block = llmBlockFactory.create(LLMBlockConfiguration.builder() + .name("Global consumer") + .llmDescriptor(llmBrick) + .prompt("Implement from ${{global.requirements}}") + .build()); + + assertTrue(block.getInputs().isEmpty()); + + FlowData flow = FlowData.builder() + .block(block) + .globalInput(IODescriptor.input("requirements", IOType.TEXT, false, null)) + .build(); + + ExecutionObject execObject = executionsService.createExecution("Global input flow", flow); + + assertEquals(ExecutionStatus.CREATED, execObject.getContext().getStatus()); + assertEquals(List.of("requirements"), execObject.getMissingGlobalInputKeys()); + assertEquals(ExecutionVariableKind.TEXT, + execObject.getContext().getExecutionVariableDescriptor("requirements").getKind()); + + execObject = executionsService.setExecutionVariable(execObject.getId(), "requirements", "Build a CRM app"); + + assertEquals(ExecutionStatus.READY, execObject.getContext().getStatus()); + assertTrue(execObject.getMissingGlobalInputKeys().isEmpty()); + assertEquals("Build a CRM app", execObject.getContext().getExecutionVariables().get("requirements")); + } + @Test public void createExecutionAndSetInput() { ExecutionObject eo = createExecutionAndSetInputInternally();