From f66bdeee1549719db666462aa6a14d2eab5366aa Mon Sep 17 00:00:00 2001 From: Lucio Lelii Date: Tue, 5 May 2026 19:58:47 +0200 Subject: [PATCH] Add execution metrics and subflow performance telemetry --- pom.xml | 4 ++ .../manager/executions/ExecutionsService.java | 45 +++++++++++- .../containers/IteratorContainerExecutor.java | 70 ++++++++++++++----- .../containers/LoopContainerExecutor.java | 70 ++++++++++++++----- .../resources/application-prod.properties | 2 +- 5 files changed, 155 insertions(+), 36 deletions(-) diff --git a/pom.xml b/pom.xml index 74bc317..b86cd38 100644 --- a/pom.xml +++ b/pom.xml @@ -37,6 +37,10 @@ org.springframework.boot spring-boot-starter-actuator + + io.micrometer + micrometer-registry-prometheus + org.springframework.boot 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 ff0b4ca..f37fea9 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 @@ -8,6 +8,8 @@ import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import java.util.stream.Collectors; +import jakarta.annotation.PostConstruct; + import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; import org.springframework.data.domain.Page; @@ -18,6 +20,10 @@ import org.springframework.transaction.annotation.Transactional; import org.springframework.web.server.ResponseStatusException; import org.springframework.http.HttpStatus; +import io.micrometer.core.instrument.Counter; +import io.micrometer.core.instrument.Gauge; +import io.micrometer.core.instrument.MeterRegistry; + import it.cnr.isti.workflow.manager.blocks.Block; import it.cnr.isti.workflow.manager.blocks.configurations.ConditionalBlockConfiguration; import it.cnr.isti.workflow.manager.blocks.configurations.ChatInteractionBlockConfiguration; @@ -43,6 +49,8 @@ import it.cnr.isti.workflow.manager.mcp.MCPSharedSessionRegistry; @Service public class ExecutionsService { + private static final org.slf4j.Logger logger = org.slf4j.LoggerFactory.getLogger(ExecutionsService.class); + private final Map executions = new ConcurrentHashMap<>(); private final ConcurrentMap lastAccessByExecutionId = new ConcurrentHashMap<>(); @@ -64,6 +72,22 @@ public class ExecutionsService { @Autowired MCPAgentService mcpAgentService; + @Autowired(required = false) + MeterRegistry meterRegistry; + + @PostConstruct + void registerMetrics() { + if (meterRegistry == null) { + return; + } + Gauge.builder("workflow.executions.cache.size", executions, Map::size) + .description("Number of execution objects currently cached in memory") + .register(meterRegistry); + Gauge.builder("workflow.executions.cache.last_access.size", lastAccessByExecutionId, Map::size) + .description("Number of last-access entries tracked for execution cache") + .register(meterRegistry); + } + @Transactional public ExecutionObject createExecution(String executionName, FlowData flow) { return createExecution(executionName, flow, null); @@ -113,11 +137,14 @@ public class ExecutionsService { public ExecutionObject getExecution(String id) { ExecutionObject toReturn = executions.get(id); if (toReturn == null) { + incrementCounter("workflow.executions.cache.miss"); ExecutionEntity entity = executionRepository.findById(id) .orElseThrow(() -> new ResponseStatusException(HttpStatus.NOT_FOUND, "Execution with id " + id + " not found")); toReturn = rebuildExecution(entity); executions.put(id, toReturn); + } else { + incrementCounter("workflow.executions.cache.hit"); } touchExecution(id); evictIfNeeded(); @@ -132,6 +159,7 @@ public class ExecutionsService { ExecutionObject inMemory = executions.get(id); if (inMemory != null) { + incrementCounter("workflow.executions.cache.hit"); if (!owner.equals(inMemory.getOwner())) { throw new ResponseStatusException(HttpStatus.FORBIDDEN, "Execution with id " + id + " is not accessible"); @@ -140,6 +168,8 @@ public class ExecutionsService { return inMemory; } + incrementCounter("workflow.executions.cache.miss"); + ExecutionEntity entity = executionRepository.findByIdAndOwner(id, owner) .orElseThrow(() -> new ResponseStatusException(HttpStatus.NOT_FOUND, "Execution with id " + id + " not found")); @@ -556,7 +586,7 @@ public class ExecutionsService { if (executions.size() <= maxInMemoryExecutions) { break; } - removeFromMemory(executionId); + removeFromMemory(executionId, "capacity"); } } @@ -572,19 +602,28 @@ public class ExecutionsService { return execution != null && execution.getContext().getStatus().isFinalState(); }) .toList(); - expiredIds.forEach(this::removeFromMemory); + expiredIds.forEach(id -> removeFromMemory(id, "ttl")); } - private void removeFromMemory(String executionId) { + private void removeFromMemory(String executionId, String reason) { ExecutionObject removed = executions.remove(executionId); lastAccessByExecutionId.remove(executionId); if (removed == null) { return; } + incrementCounter("workflow.executions.cache.eviction"); + logger.debug("Evicted execution {} from in-memory cache (reason={})", executionId, reason); cleanupManagedResourcesIfFinal(removed); removed.shutdown(); } + private void incrementCounter(String metricName) { + if (meterRegistry == null) { + return; + } + Counter.builder(metricName).register(meterRegistry).increment(); + } + private static class RequirementAccumulator { private final String key; private final String provider; diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/IteratorContainerExecutor.java b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/IteratorContainerExecutor.java index 639aded..943422e 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/IteratorContainerExecutor.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/IteratorContainerExecutor.java @@ -5,8 +5,12 @@ import java.util.LinkedHashMap; import java.util.List; import java.util.Map; +import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Component; +import io.micrometer.core.instrument.MeterRegistry; +import io.micrometer.core.instrument.Timer; + 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; @@ -25,10 +29,15 @@ import it.cnr.isti.workflow.manager.executions.steps.Input; @Component public class IteratorContainerExecutor implements ContainerExecutor { + private static final org.slf4j.Logger logger = org.slf4j.LoggerFactory.getLogger(IteratorContainerExecutor.class); private static final long INNER_EXECUTION_WAIT_TIMEOUT_MS = 60_000L; + private static final long SLOW_SUBFLOW_THRESHOLD_MS = 3_000L; private final ExecutionsService executionsService; + @Autowired(required = false) + MeterRegistry meterRegistry; + public IteratorContainerExecutor(ExecutionsService executionsService) { this.executionsService = executionsService; } @@ -129,27 +138,56 @@ public class IteratorContainerExecutor implements ContainerExecutor= SLOW_SUBFLOW_THRESHOLD_MS) { + logger.warn("Iterator container subflow was slow: executionId={}, status={}, elapsedMs={}", + innerExecution.getId(), finalStatus, elapsedMs); } } + } - if (innerExecution.getContext().getStatus() == ExecutionStatus.WAITING) { - throw new IllegalStateException( - "IteratorContainer subflow reached WAITING state. Interactive blocks inside containers are not supported yet"); + private void recordSubflowDuration(String status, long elapsedNanos) { + if (meterRegistry == null) { + return; } - if (innerExecution.getContext().getStatus() == ExecutionStatus.ERROR) { - throw new IllegalStateException("IteratorContainer subflow failed: " + innerExecution.getContext().getErrors()); - } - if (innerExecution.getContext().getStatus() != ExecutionStatus.SUCCESS) { - throw new IllegalStateException( - "IteratorContainer subflow ended in unexpected status: " + innerExecution.getContext().getStatus()); - } - return innerExecution; + Timer.builder("workflow.container.subflow.duration") + .description("Duration of container subflow executions") + .tag("containerType", IteratorContainerType.TYPE) + .tag("status", status) + .register(meterRegistry) + .record(elapsedNanos, java.util.concurrent.TimeUnit.NANOSECONDS); } private void forwardInnerEvents(ExecutionObject innerExecution, ExecutionEventLogger eventLogger, Integer iterationIndex, 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 b2d4f8d..c720ef5 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 @@ -8,6 +8,7 @@ import java.util.regex.Matcher; import java.util.regex.Pattern; import java.util.stream.Collectors; +import org.springframework.beans.factory.annotation.Autowired; import org.springframework.context.expression.MapAccessor; import org.springframework.expression.ExpressionParser; import org.springframework.expression.spel.standard.SpelExpressionParser; @@ -15,6 +16,9 @@ import org.springframework.expression.spel.support.StandardEvaluationContext; import org.springframework.stereotype.Component; import org.springframework.util.StringUtils; +import io.micrometer.core.instrument.MeterRegistry; +import io.micrometer.core.instrument.Timer; + import tools.jackson.core.type.TypeReference; import tools.jackson.databind.ObjectMapper; @@ -39,8 +43,10 @@ import it.cnr.isti.workflow.manager.llms.providers.LLMProvider; @Component public class LoopContainerExecutor implements ContainerExecutor { + private static final org.slf4j.Logger logger = org.slf4j.LoggerFactory.getLogger(LoopContainerExecutor.class); private static final Pattern PLACEHOLDER_PATTERN = Pattern.compile("\\$\\{\\{(.*?)}}"); private static final long INNER_EXECUTION_WAIT_TIMEOUT_MS = 60_000L; + private static final long SLOW_SUBFLOW_THRESHOLD_MS = 3_000L; private static final String LLM_SYSTEM_PROMPT = """ You are a workflow loop evaluator. Decide whether the loop should stop using the given prompt, inputs, outputs, execution variables and iteration index. @@ -59,6 +65,9 @@ public class LoopContainerExecutor implements ContainerExecutor llmProviders) { this.executionsService = executionsService; this.llmProviders = llmProviders; @@ -381,27 +390,56 @@ public class LoopContainerExecutor implements ContainerExecutor= SLOW_SUBFLOW_THRESHOLD_MS) { + logger.warn("Loop container subflow was slow: executionId={}, status={}, elapsedMs={}", + innerExecution.getId(), finalStatus, elapsedMs); } } + } - if (innerExecution.getContext().getStatus() == ExecutionStatus.WAITING) { - throw new IllegalStateException( - "LoopContainer subflow reached WAITING state. Interactive blocks inside containers are not supported yet"); + private void recordSubflowDuration(String status, long elapsedNanos) { + if (meterRegistry == null) { + return; } - if (innerExecution.getContext().getStatus() == ExecutionStatus.ERROR) { - throw new IllegalStateException("LoopContainer subflow failed: " + innerExecution.getContext().getErrors()); - } - if (innerExecution.getContext().getStatus() != ExecutionStatus.SUCCESS) { - throw new IllegalStateException( - "LoopContainer subflow ended in unexpected status: " + innerExecution.getContext().getStatus()); - } - return innerExecution; + Timer.builder("workflow.container.subflow.duration") + .description("Duration of container subflow executions") + .tag("containerType", LoopContainerType.TYPE) + .tag("status", status) + .register(meterRegistry) + .record(elapsedNanos, java.util.concurrent.TimeUnit.NANOSECONDS); } private void forwardInnerEvents(ExecutionObject innerExecution, ExecutionEventLogger eventLogger, Integer iterationIndex, diff --git a/src/main/resources/application-prod.properties b/src/main/resources/application-prod.properties index 10501f9..cb14e75 100644 --- a/src/main/resources/application-prod.properties +++ b/src/main/resources/application-prod.properties @@ -6,7 +6,7 @@ app.security.require-explicit-key=true app.auth.cookie.secure=true management.endpoint.health.show-details=when_authorized -management.endpoints.web.exposure.include=health,info +management.endpoints.web.exposure.include=health,info,metrics,prometheus logging.level.it.cnr.isti.workflow.manager=INFO logging.level.root=WARN