Add execution events for container internals
This commit is contained in:
parent
d3563480e4
commit
f89cec6e52
|
|
@ -14,6 +14,11 @@ public enum ExecutionEventType {
|
|||
STEP_FAILED,
|
||||
STEP_WAITING_FOR_INTERACTION,
|
||||
STEP_RESUMED,
|
||||
CONTAINER_SUBFLOW_STARTED,
|
||||
CONTAINER_SUBFLOW_COMPLETED,
|
||||
CONTAINER_ITERATION_STARTED,
|
||||
CONTAINER_ITERATION_COMPLETED,
|
||||
CONTAINER_CONDITION_EVALUATED,
|
||||
LLM_REQUEST,
|
||||
HTTP_REQUEST,
|
||||
MCP_SESSION_OPENED,
|
||||
|
|
|
|||
|
|
@ -43,7 +43,7 @@ public final class NodeExecutors {
|
|||
}
|
||||
if (node instanceof Container<?> container) {
|
||||
return ContainerExecutors.get(container.getType()).execute((Container) container, inputs, authorizations,
|
||||
executionVariables, executionVariableDescriptors);
|
||||
executionVariables, executionVariableDescriptors, eventLogger);
|
||||
}
|
||||
throw new IllegalStateException("No executor found for node type " + node.getClass().getName());
|
||||
}
|
||||
|
|
|
|||
|
|
@ -4,6 +4,7 @@ import java.util.List;
|
|||
import java.util.Map;
|
||||
|
||||
import it.cnr.isti.workflow.manager.containers.Container;
|
||||
import it.cnr.isti.workflow.manager.executions.ExecutionEventLogger;
|
||||
import it.cnr.isti.workflow.manager.containers.types.ContainerType;
|
||||
import it.cnr.isti.workflow.manager.executions.ExecutionVariableDescriptor;
|
||||
import it.cnr.isti.workflow.manager.executions.steps.Input;
|
||||
|
|
@ -11,7 +12,8 @@ import it.cnr.isti.workflow.manager.executions.steps.Input;
|
|||
public interface ContainerExecutor<T extends ContainerType> {
|
||||
|
||||
Map<String, Object> execute(Container<T> container, List<Input> inputs, Map<String, Object> authorizations,
|
||||
Map<String, Object> executionVariables, Map<String, ExecutionVariableDescriptor> executionVariableDescriptors);
|
||||
Map<String, Object> executionVariables, Map<String, ExecutionVariableDescriptor> executionVariableDescriptors,
|
||||
ExecutionEventLogger eventLogger);
|
||||
|
||||
Class<T> getContainerType();
|
||||
}
|
||||
|
|
|
|||
|
|
@ -11,6 +11,8 @@ import it.cnr.isti.workflow.manager.containers.Container;
|
|||
import it.cnr.isti.workflow.manager.containers.configurations.GenericContainerConfiguration;
|
||||
import it.cnr.isti.workflow.manager.containers.iresolvers.ContainerFlowInterfaceResolver;
|
||||
import it.cnr.isti.workflow.manager.containers.types.GenericContainerType;
|
||||
import it.cnr.isti.workflow.manager.executions.ExecutionEventLogger;
|
||||
import it.cnr.isti.workflow.manager.executions.ExecutionEventType;
|
||||
import it.cnr.isti.workflow.manager.executions.ExecutionObject;
|
||||
import it.cnr.isti.workflow.manager.executions.ExecutionStatus;
|
||||
import it.cnr.isti.workflow.manager.executions.ExecutionVariableDescriptor;
|
||||
|
|
@ -27,7 +29,7 @@ public class GenericContainerExecutor implements ContainerExecutor<GenericContai
|
|||
@Override
|
||||
public Map<String, Object> execute(Container<GenericContainerType> container, List<Input> inputs,
|
||||
Map<String, Object> authorizations, Map<String, Object> executionVariables,
|
||||
Map<String, ExecutionVariableDescriptor> executionVariableDescriptors) {
|
||||
Map<String, ExecutionVariableDescriptor> executionVariableDescriptors, ExecutionEventLogger eventLogger) {
|
||||
GenericContainerConfiguration configuration = (GenericContainerConfiguration) container.getSpecificConfiguration();
|
||||
if (configuration.getSubFlow() == null
|
||||
|| ((configuration.getSubFlow().getBlocks() == null || configuration.getSubFlow().getBlocks().isEmpty())
|
||||
|
|
@ -35,6 +37,11 @@ public class GenericContainerExecutor implements ContainerExecutor<GenericContai
|
|||
throw new IllegalArgumentException("GenericContainer subFlow is empty");
|
||||
}
|
||||
|
||||
if (eventLogger != null) {
|
||||
eventLogger.info(ExecutionEventType.CONTAINER_SUBFLOW_STARTED,
|
||||
"Starting container subFlow",
|
||||
Map.of("containerType", GenericContainerType.TYPE));
|
||||
}
|
||||
ExecutionObject innerExecution = executionsService.createExecution(container.getName() + " subflow", configuration.getSubFlow());
|
||||
|
||||
propagateAuthorizations(innerExecution, authorizations);
|
||||
|
|
@ -72,6 +79,11 @@ public class GenericContainerExecutor implements ContainerExecutor<GenericContai
|
|||
if (innerExecution.getContext().getStatus() != ExecutionStatus.SUCCESS) {
|
||||
throw new IllegalStateException("GenericContainer subflow ended in unexpected status: " + innerExecution.getContext().getStatus());
|
||||
}
|
||||
if (eventLogger != null) {
|
||||
eventLogger.info(ExecutionEventType.CONTAINER_SUBFLOW_COMPLETED,
|
||||
"Completed container subFlow",
|
||||
Map.of("containerType", GenericContainerType.TYPE, "status", innerExecution.getContext().getStatus().name()));
|
||||
}
|
||||
executionVariables.clear();
|
||||
executionVariables.putAll(innerExecution.getContext().getExecutionVariables());
|
||||
executionVariableDescriptors.clear();
|
||||
|
|
|
|||
|
|
@ -11,6 +11,8 @@ import it.cnr.isti.workflow.manager.containers.Container;
|
|||
import it.cnr.isti.workflow.manager.containers.configurations.IteratorContainerConfiguration;
|
||||
import it.cnr.isti.workflow.manager.containers.iresolvers.IteratorContainerInterfaceResolver;
|
||||
import it.cnr.isti.workflow.manager.containers.types.IteratorContainerType;
|
||||
import it.cnr.isti.workflow.manager.executions.ExecutionEventLogger;
|
||||
import it.cnr.isti.workflow.manager.executions.ExecutionEventType;
|
||||
import it.cnr.isti.workflow.manager.executions.ExecutionObject;
|
||||
import it.cnr.isti.workflow.manager.executions.ExecutionStatus;
|
||||
import it.cnr.isti.workflow.manager.executions.ExecutionVariableDescriptor;
|
||||
|
|
@ -30,7 +32,7 @@ public class IteratorContainerExecutor implements ContainerExecutor<IteratorCont
|
|||
@Override
|
||||
public Map<String, Object> execute(Container<IteratorContainerType> container, List<Input> inputs,
|
||||
Map<String, Object> authorizations, Map<String, Object> executionVariables,
|
||||
Map<String, ExecutionVariableDescriptor> executionVariableDescriptors) {
|
||||
Map<String, ExecutionVariableDescriptor> executionVariableDescriptors, ExecutionEventLogger eventLogger) {
|
||||
IteratorContainerConfiguration configuration = (IteratorContainerConfiguration) container.getSpecificConfiguration();
|
||||
IteratorContainerInterfaceResolver.Resolution resolution = IteratorContainerInterfaceResolver.resolvePorts(configuration);
|
||||
IteratorContainerInterfaceResolver.ResolvedInput iteratedInput = resolution.inputsByPublicName()
|
||||
|
|
@ -53,7 +55,14 @@ public class IteratorContainerExecutor implements ContainerExecutor<IteratorCont
|
|||
collectedOutputs.put(output.publicName(), new ArrayList<>());
|
||||
}
|
||||
|
||||
int iterationIndex = 0;
|
||||
for (Object iterationValue : iterationValues) {
|
||||
iterationIndex++;
|
||||
if (eventLogger != null) {
|
||||
eventLogger.info(ExecutionEventType.CONTAINER_ITERATION_STARTED,
|
||||
"Starting iterator iteration " + iterationIndex,
|
||||
Map.of("containerType", IteratorContainerType.TYPE, "iteration", iterationIndex, "value", iterationValue));
|
||||
}
|
||||
ExecutionObject innerExecution = executionsService.createExecution(container.getName() + " iteration",
|
||||
configuration.getSubFlow());
|
||||
|
||||
|
|
@ -84,6 +93,11 @@ public class IteratorContainerExecutor implements ContainerExecutor<IteratorCont
|
|||
resolvedOutput.exposedHandle().handle().io().getName()));
|
||||
collectedOutputs.get(resolvedOutput.publicName()).add(value);
|
||||
}
|
||||
if (eventLogger != null) {
|
||||
eventLogger.info(ExecutionEventType.CONTAINER_ITERATION_COMPLETED,
|
||||
"Completed iterator iteration " + iterationIndex,
|
||||
Map.of("containerType", IteratorContainerType.TYPE, "iteration", iterationIndex));
|
||||
}
|
||||
}
|
||||
|
||||
return new LinkedHashMap<>(collectedOutputs);
|
||||
|
|
|
|||
|
|
@ -21,6 +21,8 @@ import it.cnr.isti.workflow.manager.containers.Container;
|
|||
import it.cnr.isti.workflow.manager.containers.configurations.LoopContainerConfiguration;
|
||||
import it.cnr.isti.workflow.manager.containers.iresolvers.ContainerFlowInterfaceResolver;
|
||||
import it.cnr.isti.workflow.manager.containers.types.LoopContainerType;
|
||||
import it.cnr.isti.workflow.manager.executions.ExecutionEventLogger;
|
||||
import it.cnr.isti.workflow.manager.executions.ExecutionEventType;
|
||||
import it.cnr.isti.workflow.manager.executions.ExecutionObject;
|
||||
import it.cnr.isti.workflow.manager.executions.ExecutionStatus;
|
||||
import it.cnr.isti.workflow.manager.executions.ExecutionTemplateResolver;
|
||||
|
|
@ -55,7 +57,7 @@ public class LoopContainerExecutor implements ContainerExecutor<LoopContainerTyp
|
|||
@Override
|
||||
public Map<String, Object> execute(Container<LoopContainerType> container, List<Input> inputs,
|
||||
Map<String, Object> authorizations, Map<String, Object> executionVariables,
|
||||
Map<String, ExecutionVariableDescriptor> executionVariableDescriptors) {
|
||||
Map<String, ExecutionVariableDescriptor> executionVariableDescriptors, ExecutionEventLogger eventLogger) {
|
||||
LoopContainerConfiguration configuration = (LoopContainerConfiguration) container.getSpecificConfiguration();
|
||||
if (configuration.getSubFlow() == null || configuration.getSubFlow().getNodes().isEmpty()) {
|
||||
throw new IllegalArgumentException("LoopContainer subFlow is empty");
|
||||
|
|
@ -74,6 +76,11 @@ public class LoopContainerExecutor implements ContainerExecutor<LoopContainerTyp
|
|||
Map<String, Object> latestOutputs = new LinkedHashMap<>();
|
||||
|
||||
for (int iteration = 1; iteration <= configuration.getMaxIterations(); iteration++) {
|
||||
if (eventLogger != null) {
|
||||
eventLogger.info(ExecutionEventType.CONTAINER_ITERATION_STARTED,
|
||||
"Starting loop iteration " + iteration,
|
||||
Map.of("containerType", LoopContainerType.TYPE, "iteration", iteration));
|
||||
}
|
||||
ExecutionObject innerExecution = executionsService.createExecution(
|
||||
container.getName() + " iteration " + iteration,
|
||||
configuration.getSubFlow());
|
||||
|
|
@ -107,7 +114,16 @@ public class LoopContainerExecutor implements ContainerExecutor<LoopContainerTyp
|
|||
|
||||
updateInputsForNextIteration(currentInputs, latestOutputs, inputPortsByName);
|
||||
|
||||
if (shouldStop(configuration, currentInputs, latestOutputs, authorizations, executionVariables, iteration)) {
|
||||
boolean shouldStop = shouldStop(configuration, currentInputs, latestOutputs, authorizations, executionVariables, iteration);
|
||||
if (eventLogger != null) {
|
||||
eventLogger.info(ExecutionEventType.CONTAINER_CONDITION_EVALUATED,
|
||||
"Evaluated loop stop condition",
|
||||
Map.of("containerType", LoopContainerType.TYPE, "iteration", iteration, "stop", shouldStop));
|
||||
eventLogger.info(ExecutionEventType.CONTAINER_ITERATION_COMPLETED,
|
||||
shouldStop ? "Loop completed at iteration " + iteration : "Completed loop iteration " + iteration,
|
||||
Map.of("containerType", LoopContainerType.TYPE, "iteration", iteration, "stop", shouldStop));
|
||||
}
|
||||
if (shouldStop) {
|
||||
return new LinkedHashMap<>(latestOutputs);
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1048,6 +1048,10 @@ public class ExecutionTest {
|
|||
Object output = execObject.getContext().getResult().values().stream().findFirst().orElseThrow();
|
||||
assertTrue(output instanceof List<?>);
|
||||
assertEquals(List.of("Hello, Alice!", "Hello, Bob!"), output);
|
||||
assertTrue(execObject.getContext().getEvents().stream()
|
||||
.anyMatch(event -> event.getType() == ExecutionEventType.CONTAINER_ITERATION_STARTED));
|
||||
assertTrue(execObject.getContext().getEvents().stream()
|
||||
.anyMatch(event -> event.getType() == ExecutionEventType.CONTAINER_ITERATION_COMPLETED));
|
||||
}
|
||||
|
||||
@Test
|
||||
|
|
@ -1084,6 +1088,8 @@ public class ExecutionTest {
|
|||
assertEquals(ExecutionStatus.SUCCESS, execObject.getContext().getStatus());
|
||||
Object output = execObject.getContext().getResult().values().stream().findFirst().orElseThrow();
|
||||
assertEquals("Hello, Alice!", output);
|
||||
assertTrue(execObject.getContext().getEvents().stream()
|
||||
.anyMatch(event -> event.getType() == ExecutionEventType.CONTAINER_CONDITION_EVALUATED));
|
||||
}
|
||||
|
||||
@Test
|
||||
|
|
|
|||
Loading…
Reference in New Issue