Add execution metrics and subflow performance telemetry

This commit is contained in:
Lucio Lelii 2026-05-05 19:58:47 +02:00
parent 48925e19d4
commit f66bdeee15
5 changed files with 155 additions and 36 deletions

View File

@ -37,6 +37,10 @@
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-actuator</artifactId>
</dependency>
<dependency>
<groupId>io.micrometer</groupId>
<artifactId>micrometer-registry-prometheus</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>

View File

@ -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<String, ExecutionObject> executions = new ConcurrentHashMap<>();
private final ConcurrentMap<String, Long> 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;

View File

@ -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<IteratorContainerType> {
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<IteratorCont
}
private ExecutionObject startAndWait(ExecutionObject innerExecution) {
long startedAtNanos = System.nanoTime();
String finalStatus = "unknown";
innerExecution = executionsService.startExecution(innerExecution.getId());
while (innerExecution.getContext().getStatus() == ExecutionStatus.RUNNING) {
ExecutionStatus observedStatus = innerExecution.getContext()
.awaitStatusChangeWhileRunning(INNER_EXECUTION_WAIT_TIMEOUT_MS);
if (observedStatus == ExecutionStatus.RUNNING) {
throw new IllegalStateException("IteratorContainer subflow timed out while waiting for completion");
try {
while (innerExecution.getContext().getStatus() == ExecutionStatus.RUNNING) {
ExecutionStatus observedStatus = innerExecution.getContext()
.awaitStatusChangeWhileRunning(INNER_EXECUTION_WAIT_TIMEOUT_MS);
if (observedStatus == ExecutionStatus.RUNNING) {
finalStatus = "timeout";
throw new IllegalStateException("IteratorContainer subflow timed out while waiting for completion");
}
}
if (innerExecution.getContext().getStatus() == ExecutionStatus.WAITING) {
finalStatus = "waiting";
throw new IllegalStateException(
"IteratorContainer subflow reached WAITING state. Interactive blocks inside containers are not supported yet");
}
if (innerExecution.getContext().getStatus() == ExecutionStatus.ERROR) {
finalStatus = "error";
throw new IllegalStateException("IteratorContainer subflow failed: " + innerExecution.getContext().getErrors());
}
if (innerExecution.getContext().getStatus() != ExecutionStatus.SUCCESS) {
finalStatus = innerExecution.getContext().getStatus().name().toLowerCase();
throw new IllegalStateException(
"IteratorContainer subflow ended in unexpected status: " + innerExecution.getContext().getStatus());
}
finalStatus = "success";
return innerExecution;
} finally {
long elapsedNanos = System.nanoTime() - startedAtNanos;
recordSubflowDuration(finalStatus, elapsedNanos);
long elapsedMs = elapsedNanos / 1_000_000L;
if (elapsedMs >= 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,

View File

@ -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<LoopContainerType> {
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<LoopContainerTyp
private final ExpressionParser expressionParser = new SpelExpressionParser();
private final ObjectMapper objectMapper = new ObjectMapper();
@Autowired(required = false)
MeterRegistry meterRegistry;
public LoopContainerExecutor(ExecutionsService executionsService, Map<String, LLMProvider> llmProviders) {
this.executionsService = executionsService;
this.llmProviders = llmProviders;
@ -381,27 +390,56 @@ public class LoopContainerExecutor implements ContainerExecutor<LoopContainerTyp
}
private ExecutionObject startAndWait(ExecutionObject innerExecution) {
long startedAtNanos = System.nanoTime();
String finalStatus = "unknown";
innerExecution = executionsService.startExecution(innerExecution.getId());
while (innerExecution.getContext().getStatus() == ExecutionStatus.RUNNING) {
ExecutionStatus observedStatus = innerExecution.getContext()
.awaitStatusChangeWhileRunning(INNER_EXECUTION_WAIT_TIMEOUT_MS);
if (observedStatus == ExecutionStatus.RUNNING) {
throw new IllegalStateException("LoopContainer subflow timed out while waiting for completion");
try {
while (innerExecution.getContext().getStatus() == ExecutionStatus.RUNNING) {
ExecutionStatus observedStatus = innerExecution.getContext()
.awaitStatusChangeWhileRunning(INNER_EXECUTION_WAIT_TIMEOUT_MS);
if (observedStatus == ExecutionStatus.RUNNING) {
finalStatus = "timeout";
throw new IllegalStateException("LoopContainer subflow timed out while waiting for completion");
}
}
if (innerExecution.getContext().getStatus() == ExecutionStatus.WAITING) {
finalStatus = "waiting";
throw new IllegalStateException(
"LoopContainer subflow reached WAITING state. Interactive blocks inside containers are not supported yet");
}
if (innerExecution.getContext().getStatus() == ExecutionStatus.ERROR) {
finalStatus = "error";
throw new IllegalStateException("LoopContainer subflow failed: " + innerExecution.getContext().getErrors());
}
if (innerExecution.getContext().getStatus() != ExecutionStatus.SUCCESS) {
finalStatus = innerExecution.getContext().getStatus().name().toLowerCase();
throw new IllegalStateException(
"LoopContainer subflow ended in unexpected status: " + innerExecution.getContext().getStatus());
}
finalStatus = "success";
return innerExecution;
} finally {
long elapsedNanos = System.nanoTime() - startedAtNanos;
recordSubflowDuration(finalStatus, elapsedNanos);
long elapsedMs = elapsedNanos / 1_000_000L;
if (elapsedMs >= 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,

View File

@ -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