Add global flow inputs and execution gating

This commit is contained in:
Lucio Lelii 2026-03-30 17:07:13 +02:00
parent 7cbd655778
commit d3b5fcd9fd
9 changed files with 116 additions and 8 deletions

View File

@ -8,6 +8,7 @@ final class PlaceholderInputs {
static boolean isRuntimeExecutionVariable(String placeholder) {
return placeholder != null
&& (placeholder.startsWith("vars.")
|| placeholder.startsWith("global.")
|| placeholder.startsWith("context."));
}
}

View File

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

View File

@ -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<ExecutionAuthorizationRequirement> requiredAuthorizations = List.of();
List<IODescriptor> requiredGlobalInputs = List.of();
Map<String, Object> 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<Step<?>> 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<String, Object> executionVariables) {
this.context.setExecutionVariables(ExecutionRuntimeContextSupport.apply(executionVariables, this.id, this.name, this.owner,
this.creationTime));
registerGlobalInputs();
refreshInitializationStatus();
}
protected void setExecutionVariableDescriptors(Map<String, ExecutionVariableDescriptor> 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<String> 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;
};
}
}

View File

@ -36,6 +36,7 @@ public final class ExecutionTemplateResolver {
if (executionVariables != null) {
for (Map.Entry<String, Object> 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()));
}

View File

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

View File

@ -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<Dependency> dependencies;
@Singular
List<IODescriptor> globalInputs;
}

View File

@ -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<Dependency> dependencies;
@Singular
List<IODescriptor> globalInputs;
@JsonIgnore
public List<FlowNode> getNodes() {
List<Block<?>> blockList = blocks == null ? List.of() : blocks;

View File

@ -58,6 +58,9 @@ public class FlowDataValidator implements ConstraintValidator<ValidFlowStructure
List<FlowNode> nodes = flowData.getNodes();
List<Connection> connections = flowData.getConnections() == null ? List.of() : flowData.getConnections();
List<Dependency> dependencies = flowData.getDependencies() == null ? List.of() : flowData.getDependencies();
List<IODescriptor> globalInputs = flowData.getGlobalInputs() == null ? List.of() : flowData.getGlobalInputs();
validateGlobalInputs(globalInputs);
HashMap<String, FlowNode> nodesById = new HashMap<>();
for (FlowNode node : nodes) {
@ -76,6 +79,24 @@ public class FlowDataValidator implements ConstraintValidator<ValidFlowStructure
validateExclusiveRoutingBranches(nodes, connections, nodesById);
}
private void validateGlobalInputs(List<IODescriptor> globalInputs) {
Set<String> 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"));

View File

@ -182,6 +182,35 @@ public class ExecutionTest {
assertEquals(ExecutionStatus.CREATED, execObject.getContext().getStatus());
}
@Test
public void globalInputsAreRequiredOnceAndDoNotCreateNodeInputs() {
Block<LLMBlockType> 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();