diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/ContainerExecutionSupport.java b/src/main/java/it/cnr/isti/workflow/manager/executions/ContainerExecutionSupport.java new file mode 100644 index 0000000..5501d2e --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/ContainerExecutionSupport.java @@ -0,0 +1,89 @@ +package it.cnr.isti.workflow.manager.executions; + +import java.util.LinkedHashMap; +import java.util.LinkedHashSet; +import java.util.List; +import java.util.Map; +import java.util.Set; +import java.util.stream.Collectors; + +import it.cnr.isti.workflow.manager.containers.iresolvers.ContainerFlowInterfaceResolver; +import it.cnr.isti.workflow.manager.ios.IODescriptor; + +final class ContainerExecutionSupport { + + private ContainerExecutionSupport() { + } + + static Map collectExposedOutputsAsMap(ExecutionObject execution, + List exposedOutputs) { + Map outputs = new LinkedHashMap<>(); + for (var exposedHandle : exposedOutputs) { + var handle = exposedHandle.handle(); + Object value = execution.getContext().getResult().get(new FieldKey(handle.blockId(), handle.io().getName())); + if (value != null) { + outputs.put(exposedHandle.publicName(), value); + } + } + return outputs; + } + + static Set exposedOutputPublicNames( + List exposedOutputs) { + return exposedOutputs.stream() + .map(ContainerFlowInterfaceResolver.ExposedHandle::publicName) + .collect(Collectors.toCollection(LinkedHashSet::new)); + } + + static ExecutionVariableKind mapExecutionVariableKind(IODescriptor descriptor) { + return switch (descriptor.getType()) { + case TEXT -> ExecutionVariableKind.TEXT; + case BOOLEAN -> ExecutionVariableKind.BOOLEAN; + case FILE, CSV -> ExecutionVariableKind.FILE_PATH; + case JSON -> ExecutionVariableKind.JSON; + case ANY -> ExecutionVariableKind.ANY; + }; + } + + static Map iteratorState(List remainingValues, Map runtimeInputValues) { + Map state = new LinkedHashMap<>(); + state.put("remainingValues", new java.util.ArrayList<>(remainingValues)); + state.put("runtimeInputValues", new LinkedHashMap<>(runtimeInputValues)); + return state; + } + + static Map loopState(Map currentInputs, Map latestOutputs) { + Map state = new LinkedHashMap<>(); + state.put("currentInputs", new LinkedHashMap<>(currentInputs)); + state.put("latestOutputs", new LinkedHashMap<>(latestOutputs == null ? Map.of() : latestOutputs)); + return state; + } + + static Map widen(Map> accumulated) { + return new LinkedHashMap<>(accumulated); + } + + @SuppressWarnings("unchecked") + static List asObjectList(Object value) { + if (value instanceof List list) { + return new java.util.ArrayList<>((List) list); + } + return new java.util.ArrayList<>(); + } + + @SuppressWarnings("unchecked") + static Map asStringObjectMap(Object value) { + if (value instanceof Map map) { + return new LinkedHashMap<>((Map) map); + } + return new LinkedHashMap<>(); + } + + static Map> asAccumulatedOutputs(Map raw) { + Map> result = new LinkedHashMap<>(); + if (raw != null) { + raw.forEach((key, value) -> result.put(key, asObjectList(value))); + } + return result; + } +} 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 a7aa83f..9b932f8 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 @@ -725,8 +725,8 @@ public class ExecutionsService { // lets the container act as a branching node - downstream steps and BranchRejoins // can skip cleanly instead of erroring. parent.completeContainerSubflow(parentStep.getId(), - NodeExecutionResult.routed(collectExposedOutputsAsMap(child, exposedOutputs), - exposedOutputPublicNames(exposedOutputs))); + NodeExecutionResult.routed(ContainerExecutionSupport.collectExposedOutputsAsMap(child, exposedOutputs), + ContainerExecutionSupport.exposedOutputPublicNames(exposedOutputs))); return; } parent.failContainerSubflow(parentStep.getId(), @@ -753,7 +753,7 @@ public class ExecutionsService { accumulated.put(output.publicName(), new java.util.ArrayList<>()); } if (iterationValues.isEmpty()) { - return NodeExecutionResult.completed(widen(accumulated)); + return NodeExecutionResult.completed(ContainerExecutionSupport.widen(accumulated)); } ExecutionObject parent = getExecution(parentExecutionId); List remaining = new java.util.ArrayList<>(iterationValues.subList(1, iterationValues.size())); @@ -775,8 +775,8 @@ public class ExecutionsService { IteratorContainerInterfaceResolver.Resolution resolution = IteratorContainerInterfaceResolver.resolvePorts(configuration); List outputHandles = resolution.resolvedOutputs().stream() .map(IteratorContainerInterfaceResolver.ResolvedOutput::exposedHandle).toList(); - Map iterationOutputs = collectExposedOutputsAsMap(child, outputHandles); - Map> accumulated = asAccumulatedOutputs(continuation.getAccumulatedOutputs()); + Map iterationOutputs = ContainerExecutionSupport.collectExposedOutputsAsMap(child, outputHandles); + Map> accumulated = ContainerExecutionSupport.asAccumulatedOutputs(continuation.getAccumulatedOutputs()); for (var output : resolution.resolvedOutputs()) { accumulated.computeIfAbsent(output.publicName(), ignored -> new java.util.ArrayList<>()) .add(iterationOutputs.get(output.publicName())); @@ -788,12 +788,12 @@ public class ExecutionsService { Map.of("containerType", IteratorContainerType.TYPE, "iteration", completedIndex)); } Map state = continuation.getState() == null ? Map.of() : continuation.getState(); - List remaining = asObjectList(state.get("remainingValues")); - Map runtimeInputValues = asStringObjectMap(state.get("runtimeInputValues")); + List remaining = ContainerExecutionSupport.asObjectList(state.get("remainingValues")); + Map runtimeInputValues = ContainerExecutionSupport.asStringObjectMap(state.get("runtimeInputValues")); if (remaining.isEmpty()) { setContainerContinuation(parent.getId(), parentStep.getId(), null); - parent.completeContainerSubflow(parentStep.getId(), NodeExecutionResult.completed(widen(accumulated))); + parent.completeContainerSubflow(parentStep.getId(), NodeExecutionResult.completed(ContainerExecutionSupport.widen(accumulated))); return; } Object nextValue = remaining.getFirst(); @@ -839,7 +839,7 @@ public class ExecutionsService { container.getName() + " iteration " + iterationIndex, configuration.getSubFlow(), inputPortsByName, inputValues, iterationIndex, ContainerSubflowRole.MAIN, authorizations, eventLogger, IteratorContainerType.TYPE, innerBiasExecutionContext, executionContext, - iteratorState(remaining, runtimeInputValues), widen(accumulatedOutputs)); + ContainerExecutionSupport.iteratorState(remaining, runtimeInputValues), ContainerExecutionSupport.widen(accumulatedOutputs)); if (!child.getContext().getStatus().isFinalState()) { return ContainerAdvanceOutcome.suspended(child.getId()); @@ -849,7 +849,7 @@ public class ExecutionsService { + (child.getContext().getErrors().isEmpty() ? "" : ": " + child.getContext().getErrors())); } parent.setExecutionVariableDescriptors(child.getContext().getExecutionVariableDescriptors()); - Map iterationOutputs = collectExposedOutputsAsMap(child, outputHandles); + Map iterationOutputs = ContainerExecutionSupport.collectExposedOutputsAsMap(child, outputHandles); for (var output : resolution.resolvedOutputs()) { accumulatedOutputs.get(output.publicName()).add(iterationOutputs.get(output.publicName())); } @@ -859,7 +859,7 @@ public class ExecutionsService { } if (remaining.isEmpty()) { setContainerContinuation(parent.getId(), parentStepId, null); - return ContainerAdvanceOutcome.completed(widen(accumulatedOutputs)); + return ContainerAdvanceOutcome.completed(ContainerExecutionSupport.widen(accumulatedOutputs)); } iterationValue = remaining.getFirst(); remaining = new java.util.ArrayList<>(remaining.subList(1, remaining.size())); @@ -894,8 +894,8 @@ public class ExecutionsService { parent.setExecutionVariableDescriptors(completedChild.getContext().getExecutionVariableDescriptors()); Map state = continuation.getState() == null ? Map.of() : continuation.getState(); - Map currentInputs = asStringObjectMap(state.get("currentInputs")); - Map latestOutputs = asStringObjectMap(state.get("latestOutputs")); + Map currentInputs = ContainerExecutionSupport.asStringObjectMap(state.get("currentInputs")); + Map latestOutputs = ContainerExecutionSupport.asStringObjectMap(state.get("latestOutputs")); int iterationIndex = continuation.getIterationIndex() == null ? 1 : continuation.getIterationIndex(); ExecutionEventLogger eventLogger = parentStep.getEventLogger(); BiasExecutionContext innerBiasExecutionContext = BiasContainerPropagation.innerContextFor(container.getId(), @@ -908,7 +908,7 @@ public class ExecutionsService { if (completedChild.getSubflowRole() != ContainerSubflowRole.GUARD) { List mainOutputHandles = ContainerFlowInterfaceResolver .getExposedOutputs(configuration.getSubFlow()); - Map freshLatest = collectExposedOutputsAsMap(completedChild, mainOutputHandles); + Map freshLatest = ContainerExecutionSupport.collectExposedOutputsAsMap(completedChild, mainOutputHandles); outcome = runLoopFrom(parent, parentStep.getId(), container, configuration, LoopPhase.GUARD, iterationIndex, currentInputs, freshLatest, authorizations, eventLogger, innerBiasExecutionContext, executionContext); } else { @@ -918,7 +918,7 @@ public class ExecutionsService { String feedbackInput = resolveFeedbackInput(configuration, inputPortsByName); List guardOutputHandles = ContainerFlowInterfaceResolver .getExposedOutputs(configuration.getGuardSubFlow()); - Map guardResult = collectExposedOutputsAsMap(completedChild, guardOutputHandles); + Map guardResult = ContainerExecutionSupport.collectExposedOutputsAsMap(completedChild, guardOutputHandles); boolean shouldContinue = parseLoopGuardResponse( String.valueOf(requireGuardOutput(guardResult, LoopContainerConfiguration.GUARD_OUTPUT))); String feedback = String.valueOf(requireGuardOutput(guardResult, LoopContainerConfiguration.FEEDBACK_OUTPUT)); @@ -976,7 +976,7 @@ public class ExecutionsService { container.getName() + " iteration " + iterationIndex, configuration.getSubFlow(), inputPortsByName, inputs, iterationIndex, ContainerSubflowRole.MAIN, authorizations, eventLogger, LoopContainerType.TYPE, innerBiasExecutionContext, executionContext, - loopState(inputs, latestOutputs), null); + ContainerExecutionSupport.loopState(inputs, latestOutputs), null); if (!mainChild.getContext().getStatus().isFinalState()) { return ContainerAdvanceOutcome.suspended(mainChild.getId()); } @@ -985,7 +985,7 @@ public class ExecutionsService { + (mainChild.getContext().getErrors().isEmpty() ? "" : ": " + mainChild.getContext().getErrors())); } parent.setExecutionVariableDescriptors(mainChild.getContext().getExecutionVariableDescriptors()); - latestOutputs = collectExposedOutputsAsMap(mainChild, mainOutputHandles); + latestOutputs = ContainerExecutionSupport.collectExposedOutputsAsMap(mainChild, mainOutputHandles); phase = LoopPhase.GUARD; continue; } @@ -998,7 +998,7 @@ public class ExecutionsService { container.getName() + " guard iteration " + iterationIndex, configuration.getGuardSubFlow(), guardInputsByName, guardTemplateValues, iterationIndex, ContainerSubflowRole.GUARD, authorizations, eventLogger, "LoopContainerGuard", innerBiasExecutionContext, executionContext, - loopState(nextInputsIfContinuing, latestOutputs), null); + ContainerExecutionSupport.loopState(nextInputsIfContinuing, latestOutputs), null); if (!guardChild.getContext().getStatus().isFinalState()) { return ContainerAdvanceOutcome.suspended(guardChild.getId()); } @@ -1007,7 +1007,7 @@ public class ExecutionsService { + (guardChild.getContext().getErrors().isEmpty() ? "" : ": " + guardChild.getContext().getErrors())); } parent.setExecutionVariableDescriptors(guardChild.getContext().getExecutionVariableDescriptors()); - Map guardResult = collectExposedOutputsAsMap(guardChild, guardOutputHandles); + Map guardResult = ContainerExecutionSupport.collectExposedOutputsAsMap(guardChild, guardOutputHandles); boolean shouldContinue = parseLoopGuardResponse( String.valueOf(requireGuardOutput(guardResult, LoopContainerConfiguration.GUARD_OUTPUT))); String feedback = String.valueOf(requireGuardOutput(guardResult, LoopContainerConfiguration.FEEDBACK_OUTPUT)); @@ -1133,7 +1133,7 @@ public class ExecutionsService { for (IODescriptor required : child.getRequiredGlobalInputs()) { globalDescriptors.put(required.getName(), ExecutionVariableDescriptor.builder() .name(required.getName()) - .kind(mapExecutionVariableKind(required)) + .kind(ContainerExecutionSupport.mapExecutionVariableKind(required)) .value(globalInputs.get(required.getName())) .cleanupPolicy(ExecutionVariableCleanupPolicy.NONE) .description("Global flow input") @@ -1205,78 +1205,6 @@ public class ExecutionsService { return details; } - private static Map collectExposedOutputsAsMap(ExecutionObject execution, - List exposedOutputs) { - Map outputs = new LinkedHashMap<>(); - for (var exposedHandle : exposedOutputs) { - var handle = exposedHandle.handle(); - Object value = execution.getContext().getResult().get(new FieldKey(handle.blockId(), handle.io().getName())); - if (value != null) { - outputs.put(exposedHandle.publicName(), value); - } - } - return outputs; - } - - private static Set exposedOutputPublicNames( - List exposedOutputs) { - return exposedOutputs.stream() - .map(ContainerFlowInterfaceResolver.ExposedHandle::publicName) - .collect(Collectors.toCollection(LinkedHashSet::new)); - } - - private static ExecutionVariableKind mapExecutionVariableKind(IODescriptor descriptor) { - return switch (descriptor.getType()) { - case TEXT -> ExecutionVariableKind.TEXT; - case BOOLEAN -> ExecutionVariableKind.BOOLEAN; - case FILE, CSV -> ExecutionVariableKind.FILE_PATH; - case JSON -> ExecutionVariableKind.JSON; - case ANY -> ExecutionVariableKind.ANY; - }; - } - - private static Map iteratorState(List remainingValues, Map runtimeInputValues) { - Map state = new LinkedHashMap<>(); - state.put("remainingValues", new java.util.ArrayList<>(remainingValues)); - state.put("runtimeInputValues", new LinkedHashMap<>(runtimeInputValues)); - return state; - } - - private static Map loopState(Map currentInputs, Map latestOutputs) { - Map state = new LinkedHashMap<>(); - state.put("currentInputs", new LinkedHashMap<>(currentInputs)); - state.put("latestOutputs", new LinkedHashMap<>(latestOutputs == null ? Map.of() : latestOutputs)); - return state; - } - - private static Map widen(Map> accumulated) { - return new LinkedHashMap<>(accumulated); - } - - @SuppressWarnings("unchecked") - private static List asObjectList(Object value) { - if (value instanceof List list) { - return new java.util.ArrayList<>((List) list); - } - return new java.util.ArrayList<>(); - } - - @SuppressWarnings("unchecked") - private static Map asStringObjectMap(Object value) { - if (value instanceof Map map) { - return new LinkedHashMap<>((Map) map); - } - return new LinkedHashMap<>(); - } - - private static Map> asAccumulatedOutputs(Map raw) { - Map> result = new LinkedHashMap<>(); - if (raw != null) { - raw.forEach((key, value) -> result.put(key, asObjectList(value))); - } - return result; - } - private record ContainerAdvanceOutcome(boolean completed, Map outputs, String pendingChildId) { static ContainerAdvanceOutcome completed(Map outputs) { return new ContainerAdvanceOutcome(true, outputs, null);