diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionEventType.java b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionEventType.java index f4d9b91..59d7015 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionEventType.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionEventType.java @@ -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, diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/NodeExecutors.java b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/NodeExecutors.java index e8fe70a..b71bbd5 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/NodeExecutors.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/NodeExecutors.java @@ -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()); } diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/ContainerExecutor.java b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/ContainerExecutor.java index 5e929f5..6184f64 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/ContainerExecutor.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/ContainerExecutor.java @@ -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 { Map execute(Container container, List inputs, Map authorizations, - Map executionVariables, Map executionVariableDescriptors); + Map executionVariables, Map executionVariableDescriptors, + ExecutionEventLogger eventLogger); Class getContainerType(); } diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/GenericContainerExecutor.java b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/GenericContainerExecutor.java index 6d0c755..5981d7e 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/GenericContainerExecutor.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/GenericContainerExecutor.java @@ -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 execute(Container container, List inputs, Map authorizations, Map executionVariables, - Map executionVariableDescriptors) { + Map 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 execute(Container container, List inputs, Map authorizations, Map executionVariables, - Map executionVariableDescriptors) { + Map 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()); } + 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(collectedOutputs); diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/LoopContainerExecutor.java b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/LoopContainerExecutor.java index 41dc988..4224b3e 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/LoopContainerExecutor.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/LoopContainerExecutor.java @@ -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 execute(Container container, List inputs, Map authorizations, Map executionVariables, - Map executionVariableDescriptors) { + Map 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 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(latestOutputs); } } 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 0066e75..0e5ee61 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 @@ -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