refactor(executions): extract container data-shape helpers into ContainerExecutionSupport
Cluster K from a structural analysis of ExecutionsService.java (1862 lines, second-largest file after FlowAssistantService): 9 fully field-free static functions for shaping container/iterator/loop execution data (collectExposedOutputsAsMap, exposedOutputPublicNames, mapExecutionVariableKind, iteratorState, loopState, widen, asObjectList, asStringObjectMap, asAccumulatedOutputs). ContainerAdvanceOutcome stays in ExecutionsService (only used by runIteratorIterations/runLoopFrom/toNodeExecutionResult, which remain there). 1862 -> 1790 lines. Behavior-preserving: pure extraction, no logic changes. Verified with `mvn test`: 462 tests, 0 failures, 0 errors. Co-Authored-By: Claude Haiku 4.5 <noreply@anthropic.com>
This commit is contained in:
parent
9c7e0233a2
commit
4cabbe90d3
|
|
@ -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<String, Object> collectExposedOutputsAsMap(ExecutionObject execution,
|
||||
List<ContainerFlowInterfaceResolver.ExposedHandle> exposedOutputs) {
|
||||
Map<String, Object> 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<String> exposedOutputPublicNames(
|
||||
List<ContainerFlowInterfaceResolver.ExposedHandle> 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<String, Object> iteratorState(List<Object> remainingValues, Map<String, Object> runtimeInputValues) {
|
||||
Map<String, Object> state = new LinkedHashMap<>();
|
||||
state.put("remainingValues", new java.util.ArrayList<>(remainingValues));
|
||||
state.put("runtimeInputValues", new LinkedHashMap<>(runtimeInputValues));
|
||||
return state;
|
||||
}
|
||||
|
||||
static Map<String, Object> loopState(Map<String, Object> currentInputs, Map<String, Object> latestOutputs) {
|
||||
Map<String, Object> state = new LinkedHashMap<>();
|
||||
state.put("currentInputs", new LinkedHashMap<>(currentInputs));
|
||||
state.put("latestOutputs", new LinkedHashMap<>(latestOutputs == null ? Map.of() : latestOutputs));
|
||||
return state;
|
||||
}
|
||||
|
||||
static Map<String, Object> widen(Map<String, List<Object>> accumulated) {
|
||||
return new LinkedHashMap<>(accumulated);
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
static List<Object> asObjectList(Object value) {
|
||||
if (value instanceof List<?> list) {
|
||||
return new java.util.ArrayList<>((List<Object>) list);
|
||||
}
|
||||
return new java.util.ArrayList<>();
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
static Map<String, Object> asStringObjectMap(Object value) {
|
||||
if (value instanceof Map<?, ?> map) {
|
||||
return new LinkedHashMap<>((Map<String, Object>) map);
|
||||
}
|
||||
return new LinkedHashMap<>();
|
||||
}
|
||||
|
||||
static Map<String, List<Object>> asAccumulatedOutputs(Map<String, Object> raw) {
|
||||
Map<String, List<Object>> result = new LinkedHashMap<>();
|
||||
if (raw != null) {
|
||||
raw.forEach((key, value) -> result.put(key, asObjectList(value)));
|
||||
}
|
||||
return result;
|
||||
}
|
||||
}
|
||||
|
|
@ -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<Object> remaining = new java.util.ArrayList<>(iterationValues.subList(1, iterationValues.size()));
|
||||
|
|
@ -775,8 +775,8 @@ public class ExecutionsService {
|
|||
IteratorContainerInterfaceResolver.Resolution resolution = IteratorContainerInterfaceResolver.resolvePorts(configuration);
|
||||
List<ContainerFlowInterfaceResolver.ExposedHandle> outputHandles = resolution.resolvedOutputs().stream()
|
||||
.map(IteratorContainerInterfaceResolver.ResolvedOutput::exposedHandle).toList();
|
||||
Map<String, Object> iterationOutputs = collectExposedOutputsAsMap(child, outputHandles);
|
||||
Map<String, List<Object>> accumulated = asAccumulatedOutputs(continuation.getAccumulatedOutputs());
|
||||
Map<String, Object> iterationOutputs = ContainerExecutionSupport.collectExposedOutputsAsMap(child, outputHandles);
|
||||
Map<String, List<Object>> 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<String, Object> state = continuation.getState() == null ? Map.of() : continuation.getState();
|
||||
List<Object> remaining = asObjectList(state.get("remainingValues"));
|
||||
Map<String, Object> runtimeInputValues = asStringObjectMap(state.get("runtimeInputValues"));
|
||||
List<Object> remaining = ContainerExecutionSupport.asObjectList(state.get("remainingValues"));
|
||||
Map<String, Object> 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<String, Object> iterationOutputs = collectExposedOutputsAsMap(child, outputHandles);
|
||||
Map<String, Object> 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<String, Object> state = continuation.getState() == null ? Map.of() : continuation.getState();
|
||||
Map<String, Object> currentInputs = asStringObjectMap(state.get("currentInputs"));
|
||||
Map<String, Object> latestOutputs = asStringObjectMap(state.get("latestOutputs"));
|
||||
Map<String, Object> currentInputs = ContainerExecutionSupport.asStringObjectMap(state.get("currentInputs"));
|
||||
Map<String, Object> 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<ContainerFlowInterfaceResolver.ExposedHandle> mainOutputHandles = ContainerFlowInterfaceResolver
|
||||
.getExposedOutputs(configuration.getSubFlow());
|
||||
Map<String, Object> freshLatest = collectExposedOutputsAsMap(completedChild, mainOutputHandles);
|
||||
Map<String, Object> 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<ContainerFlowInterfaceResolver.ExposedHandle> guardOutputHandles = ContainerFlowInterfaceResolver
|
||||
.getExposedOutputs(configuration.getGuardSubFlow());
|
||||
Map<String, Object> guardResult = collectExposedOutputsAsMap(completedChild, guardOutputHandles);
|
||||
Map<String, Object> 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<String, Object> guardResult = collectExposedOutputsAsMap(guardChild, guardOutputHandles);
|
||||
Map<String, Object> 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<String, Object> collectExposedOutputsAsMap(ExecutionObject execution,
|
||||
List<ContainerFlowInterfaceResolver.ExposedHandle> exposedOutputs) {
|
||||
Map<String, Object> 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<String> exposedOutputPublicNames(
|
||||
List<ContainerFlowInterfaceResolver.ExposedHandle> 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<String, Object> iteratorState(List<Object> remainingValues, Map<String, Object> runtimeInputValues) {
|
||||
Map<String, Object> state = new LinkedHashMap<>();
|
||||
state.put("remainingValues", new java.util.ArrayList<>(remainingValues));
|
||||
state.put("runtimeInputValues", new LinkedHashMap<>(runtimeInputValues));
|
||||
return state;
|
||||
}
|
||||
|
||||
private static Map<String, Object> loopState(Map<String, Object> currentInputs, Map<String, Object> latestOutputs) {
|
||||
Map<String, Object> state = new LinkedHashMap<>();
|
||||
state.put("currentInputs", new LinkedHashMap<>(currentInputs));
|
||||
state.put("latestOutputs", new LinkedHashMap<>(latestOutputs == null ? Map.of() : latestOutputs));
|
||||
return state;
|
||||
}
|
||||
|
||||
private static Map<String, Object> widen(Map<String, List<Object>> accumulated) {
|
||||
return new LinkedHashMap<>(accumulated);
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
private static List<Object> asObjectList(Object value) {
|
||||
if (value instanceof List<?> list) {
|
||||
return new java.util.ArrayList<>((List<Object>) list);
|
||||
}
|
||||
return new java.util.ArrayList<>();
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
private static Map<String, Object> asStringObjectMap(Object value) {
|
||||
if (value instanceof Map<?, ?> map) {
|
||||
return new LinkedHashMap<>((Map<String, Object>) map);
|
||||
}
|
||||
return new LinkedHashMap<>();
|
||||
}
|
||||
|
||||
private static Map<String, List<Object>> asAccumulatedOutputs(Map<String, Object> raw) {
|
||||
Map<String, List<Object>> result = new LinkedHashMap<>();
|
||||
if (raw != null) {
|
||||
raw.forEach((key, value) -> result.put(key, asObjectList(value)));
|
||||
}
|
||||
return result;
|
||||
}
|
||||
|
||||
private record ContainerAdvanceOutcome(boolean completed, Map<String, Object> outputs, String pendingChildId) {
|
||||
static ContainerAdvanceOutcome completed(Map<String, Object> outputs) {
|
||||
return new ContainerAdvanceOutcome(true, outputs, null);
|
||||
|
|
|
|||
Loading…
Reference in New Issue