diff --git a/pom.xml b/pom.xml
index 0fc1972..74bc317 100644
--- a/pom.xml
+++ b/pom.xml
@@ -97,6 +97,10 @@
org.postgresql
postgresql
+
+ org.flywaydb
+ flyway-core
+
org.springframework.boot
spring-boot-starter-web
@@ -168,6 +172,10 @@
org.apache.maven.plugins
maven-compiler-plugin
+ ${java.version}
+
+ --enable-preview
+
org.projectlombok
@@ -192,12 +200,30 @@
org.apache.maven.plugins
maven-surefire-plugin
- -javaagent:${settings.localRepository}/net/bytebuddy/byte-buddy-agent/${byte-buddy.version}/byte-buddy-agent-${byte-buddy.version}.jar
+ --enable-preview ${argLine} -javaagent:${settings.localRepository}/net/bytebuddy/byte-buddy-agent/${byte-buddy.version}/byte-buddy-agent-${byte-buddy.version}.jar
subclass
+
+ org.jacoco
+ jacoco-maven-plugin
+
+
+
+ prepare-agent
+
+
+
+ report
+ test
+
+ report
+
+
+
+
diff --git a/src/main/java/it/cnr/isti/workflow/manager/auth/config/JwtUtil.java b/src/main/java/it/cnr/isti/workflow/manager/auth/config/JwtUtil.java
index b71eac7..176dc10 100644
--- a/src/main/java/it/cnr/isti/workflow/manager/auth/config/JwtUtil.java
+++ b/src/main/java/it/cnr/isti/workflow/manager/auth/config/JwtUtil.java
@@ -17,7 +17,12 @@ public class JwtUtil {
SecretKey secretKey;
- public JwtUtil( @Value("${app.security.key}") String secretKeyString) {
+ public JwtUtil(@Value("${app.security.key}") String secretKeyString,
+ @Value("${app.security.default-key:}") String defaultKey,
+ @Value("${app.security.require-explicit-key:false}") boolean requireExplicitKey) {
+ if (requireExplicitKey && secretKeyString.equals(defaultKey)) {
+ throw new IllegalStateException("WFEDITOR_SECRET_KEY must be explicitly configured in this environment");
+ }
this.secretKey = Keys.hmacShaKeyFor(secretKeyString.getBytes());
}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/controllers/ExecutionsController.java b/src/main/java/it/cnr/isti/workflow/manager/controllers/ExecutionsController.java
index e6dc9bd..8aae331 100644
--- a/src/main/java/it/cnr/isti/workflow/manager/controllers/ExecutionsController.java
+++ b/src/main/java/it/cnr/isti/workflow/manager/controllers/ExecutionsController.java
@@ -355,11 +355,7 @@ public class ExecutionsController {
if (userDetails == null) {
throw new ResponseStatusException(HttpStatus.FORBIDDEN, "Authentication required");
}
- ExecutionObject execution = executionService.getExecution(id);
- if (!userDetails.getUsername().equals(execution.getOwner())) {
- throw new ResponseStatusException(HttpStatus.FORBIDDEN, "Execution with id " + id + " is not accessible");
- }
- return execution;
+ return executionService.getExecutionByOwner(id, userDetails.getUsername());
}
}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionContext.java b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionContext.java
index 32cbd92..768dcd6 100644
--- a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionContext.java
+++ b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionContext.java
@@ -71,6 +71,9 @@ public class ExecutionContext implements ExecutionListener {
@JsonIgnore
Runnable errorStateListener;
+ @JsonIgnore
+ private final Object statusMonitor = new Object();
+
boolean interactionSimulationEnabled = false;
LLMDescriptor interactionSimulationDescriptor;
@@ -144,6 +147,9 @@ public class ExecutionContext implements ExecutionListener {
.details(Map.of("status", status.name()))
.build());
}
+ synchronized (statusMonitor) {
+ statusMonitor.notifyAll();
+ }
notifyStateChanged();
}
@@ -614,6 +620,28 @@ public class ExecutionContext implements ExecutionListener {
if (this.stateChangeListener != null) {
this.stateChangeListener.run();
}
+ synchronized (statusMonitor) {
+ statusMonitor.notifyAll();
+ }
+ }
+
+ public ExecutionStatus awaitStatusChangeWhileRunning(long timeoutMs) {
+ long deadline = timeoutMs <= 0 ? Long.MAX_VALUE : System.currentTimeMillis() + timeoutMs;
+ synchronized (statusMonitor) {
+ while (this.status == ExecutionStatus.RUNNING) {
+ long remaining = deadline - System.currentTimeMillis();
+ if (remaining <= 0) {
+ break;
+ }
+ try {
+ statusMonitor.wait(remaining);
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ break;
+ }
+ }
+ return this.status;
+ }
}
protected void setRuntimeContextValues(Map runtimeContextValues) {
diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionObject.java b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionObject.java
index 31c35e3..e8cb358 100644
--- a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionObject.java
+++ b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionObject.java
@@ -35,6 +35,8 @@ import lombok.NoArgsConstructor;
@NoArgsConstructor
public class ExecutionObject {
+ private static final int MIN_EXECUTOR_THREADS = 1;
+
String id = UUID.randomUUID().toString();
ExecutionContext context;
@@ -80,7 +82,7 @@ public class ExecutionObject {
List> steps = getStepsFromFlow(flow);
- this.executorService = Executors.newFixedThreadPool(steps.size());
+ this.executorService = createExecutorService(steps.size());
this.context = new ExecutionContext(steps.stream().collect(Collectors.toMap(Step::getId, Function.identity())));
ensureRuntimeContextVariables();
@@ -238,6 +240,9 @@ public class ExecutionObject {
step.setSimulated(simulateInteractions);
}
});
+ if (this.executorService == null || this.executorService.isShutdown() || this.executorService.isTerminated()) {
+ this.executorService = createExecutorService(this.context.getSteps().size());
+ }
if (this.context.getStatus() == ExecutionStatus.READY){
this.context.start(executorService);
} else
@@ -265,12 +270,14 @@ public class ExecutionObject {
this.executorService.shutdownNow();
Thread.currentThread().interrupt();
}
+ this.executorService = null;
}
}
protected void cancel() {
if (this.executorService != null) {
this.executorService.shutdownNow();
+ this.executorService = null;
}
this.context.cancel();
}
@@ -379,9 +386,15 @@ public class ExecutionObject {
protected void abortOnError() {
if (this.executorService != null) {
this.executorService.shutdownNow();
+ this.executorService = null;
}
}
+ private ExecutorService createExecutorService(int requestedThreads) {
+ int threadCount = Math.max(MIN_EXECUTOR_THREADS, requestedThreads);
+ return Executors.newFixedThreadPool(threadCount);
+ }
+
private void ensureRuntimeContextVariables() {
Map runtimeContext = ExecutionRuntimeContextSupport.contextValues(
this.id,
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 fe04b0e..ff0b4ca 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
@@ -5,11 +5,16 @@ import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.ConcurrentMap;
+import java.util.stream.Collectors;
import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.beans.factory.annotation.Value;
import org.springframework.data.domain.Page;
import org.springframework.data.domain.Pageable;
+import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Service;
+import org.springframework.transaction.annotation.Transactional;
import org.springframework.web.server.ResponseStatusException;
import org.springframework.http.HttpStatus;
@@ -39,6 +44,13 @@ import it.cnr.isti.workflow.manager.mcp.MCPSharedSessionRegistry;
public class ExecutionsService {
private final Map executions = new ConcurrentHashMap<>();
+ private final ConcurrentMap lastAccessByExecutionId = new ConcurrentHashMap<>();
+
+ @Value("${app.executions.cache.max-size:1000}")
+ private int maxInMemoryExecutions;
+
+ @Value("${app.executions.cache.final-ttl-ms:1800000}")
+ private long finalStateTtlMs;
@Autowired
FlowExecutionValidator flowExecutionValidator;
@@ -52,10 +64,12 @@ public class ExecutionsService {
@Autowired
MCPAgentService mcpAgentService;
+ @Transactional
public ExecutionObject createExecution(String executionName, FlowData flow) {
return createExecution(executionName, flow, null);
}
+ @Transactional
public ExecutionObject createExecution(String executionName, FlowData flow, String owner) {
flowExecutionValidator.validate(flow);
List requiredAuthorizations = resolveRequiredAuthorizations(flow);
@@ -67,10 +81,13 @@ public class ExecutionsService {
.build();
attachPersistence(execObject);
executions.put(execObject.getId(), execObject);
+ touchExecution(execObject.getId());
+ evictIfNeeded();
persist(execObject);
return execObject;
}
+ @Transactional
public ExecutionObject createExecution(Flow flow) {
FlowData.FlowDataBuilder flowDataBuilder = FlowData.builder();
if (flow.getBlocks() != null) {
@@ -92,6 +109,7 @@ public class ExecutionsService {
return createExecution(flow.getName(), flowData, null);
}
+ @Transactional(readOnly = true)
public ExecutionObject getExecution(String id) {
ExecutionObject toReturn = executions.get(id);
if (toReturn == null) {
@@ -101,20 +119,55 @@ public class ExecutionsService {
toReturn = rebuildExecution(entity);
executions.put(id, toReturn);
}
+ touchExecution(id);
+ evictIfNeeded();
return toReturn;
}
+ @Transactional(readOnly = true)
+ public ExecutionObject getExecutionByOwner(String id, String owner) {
+ if (owner == null || owner.isBlank()) {
+ throw new ResponseStatusException(HttpStatus.FORBIDDEN, "Authentication required");
+ }
+
+ ExecutionObject inMemory = executions.get(id);
+ if (inMemory != null) {
+ if (!owner.equals(inMemory.getOwner())) {
+ throw new ResponseStatusException(HttpStatus.FORBIDDEN,
+ "Execution with id " + id + " is not accessible");
+ }
+ touchExecution(id);
+ return inMemory;
+ }
+
+ ExecutionEntity entity = executionRepository.findByIdAndOwner(id, owner)
+ .orElseThrow(() -> new ResponseStatusException(HttpStatus.NOT_FOUND,
+ "Execution with id " + id + " not found"));
+ ExecutionObject rebuilt = rebuildExecution(entity);
+ executions.put(id, rebuilt);
+ touchExecution(id);
+ evictIfNeeded();
+ return rebuilt;
+ }
+
public List getAllExecutions() {
- return executionRepository.findAll().stream()
- .map(entity -> executions.computeIfAbsent(entity.getId(), ignored -> rebuildExecution(entity)))
+ List loadedExecutions = executionRepository.findAll().stream()
+ .map(entity -> executions.computeIfAbsent(entity.getId(), ignored -> rebuildExecution(entity)))
.toList();
+ loadedExecutions.forEach(execution -> touchExecution(execution.getId()));
+ evictIfNeeded();
+ return loadedExecutions;
}
public Page getExecutionsByOwner(String owner, Pageable pageable) {
- return executionRepository.findByOwner(owner, pageable)
- .map(entity -> executions.computeIfAbsent(entity.getId(), ignored -> rebuildExecution(entity)));
+ Page page = executionRepository.findByOwner(owner, pageable)
+ .map(entity -> executions.computeIfAbsent(entity.getId(), ignored -> rebuildExecution(entity)));
+ page.forEach(execution -> touchExecution(execution.getId()));
+ evictIfNeeded();
+ return page;
}
+ @Transactional
public void removeExecution(String id) {
ExecutionObject execution = executions.get(id);
if (execution == null && executionRepository.existsById(id)) {
@@ -127,15 +180,18 @@ public class ExecutionsService {
throw new IllegalStateException("Execution with id " + id + " is still running");
execution.shutdown();
executions.remove(id);
+ lastAccessByExecutionId.remove(id);
executionRepository.deleteById(id);
}
+ @Transactional
public ExecutionObject cancelExecution(String id) {
ExecutionObject eo = getExecution(id);
if (eo.getContext().getStatus().isFinalState()) {
return eo;
}
eo.cancel();
+ touchExecution(id);
return eo;
}
@@ -237,12 +293,15 @@ public class ExecutionsService {
return getExecution(executionId).removeExecutionVariable(key);
}
+ @Transactional
public ExecutionObject startExecution(String id) {
ExecutionObject eo = getExecution(id);
eo.start();
+ touchExecution(id);
return eo;
}
+ @Transactional
public ExecutionObject startSimulationExecution(String id, LLMDescriptor simulatorDescriptor) {
ExecutionObject eo = getExecution(id);
if (!eo.isSimulationAvailable()) {
@@ -255,9 +314,11 @@ public class ExecutionsService {
}
eo.setInteractionSimulationDescriptor(simulatorDescriptor);
eo.startSimulation();
+ touchExecution(id);
return eo;
}
+ @Transactional
public ExecutionObject resumeExecution(String id) {
ExecutionObject eo = getExecution(id);
if (eo.getContext().getStatus().isFinalState()) {
@@ -270,11 +331,20 @@ public class ExecutionsService {
if (eo.getContext().getStatus() == ExecutionStatus.READY) {
eo.start();
}
+ touchExecution(id);
return eo;
}
public void clearInMemoryExecutions() {
+ executions.values().forEach(ExecutionObject::shutdown);
executions.clear();
+ lastAccessByExecutionId.clear();
+ }
+
+ @Scheduled(fixedDelayString = "${app.executions.cache.cleanup-interval-ms:60000}")
+ void cleanupInMemoryExecutions() {
+ evictFinalStateByTtl(System.currentTimeMillis());
+ evictIfNeeded();
}
private List resolveRequiredAuthorizations(FlowData flow) {
@@ -420,6 +490,7 @@ public class ExecutionsService {
.flow(executionObject.getFlow())
.snapshot(executionObject.snapshot())
.build());
+ touchExecution(executionObject.getId());
}
private ExecutionObject rebuildExecution(ExecutionEntity entity) {
@@ -463,6 +534,57 @@ public class ExecutionsService {
cleanedDescriptors.forEach(descriptor -> descriptor.setCleanupPolicy(ExecutionVariableCleanupPolicy.NONE));
}
+ private void touchExecution(String executionId) {
+ if (executionId != null) {
+ lastAccessByExecutionId.put(executionId, System.currentTimeMillis());
+ }
+ }
+
+ private void evictIfNeeded() {
+ if (executions.size() <= maxInMemoryExecutions) {
+ return;
+ }
+ List evictionCandidates = lastAccessByExecutionId.entrySet().stream()
+ .sorted(Map.Entry.comparingByValue())
+ .map(Map.Entry::getKey)
+ .filter(id -> {
+ ExecutionObject execution = executions.get(id);
+ return execution != null && execution.getContext().getStatus().isFinalState();
+ })
+ .collect(Collectors.toList());
+ for (String executionId : evictionCandidates) {
+ if (executions.size() <= maxInMemoryExecutions) {
+ break;
+ }
+ removeFromMemory(executionId);
+ }
+ }
+
+ private void evictFinalStateByTtl(long now) {
+ if (finalStateTtlMs <= 0) {
+ return;
+ }
+ List expiredIds = lastAccessByExecutionId.entrySet().stream()
+ .filter(entry -> now - entry.getValue() >= finalStateTtlMs)
+ .map(Map.Entry::getKey)
+ .filter(id -> {
+ ExecutionObject execution = executions.get(id);
+ return execution != null && execution.getContext().getStatus().isFinalState();
+ })
+ .toList();
+ expiredIds.forEach(this::removeFromMemory);
+ }
+
+ private void removeFromMemory(String executionId) {
+ ExecutionObject removed = executions.remove(executionId);
+ lastAccessByExecutionId.remove(executionId);
+ if (removed == null) {
+ return;
+ }
+ cleanupManagedResourcesIfFinal(removed);
+ removed.shutdown();
+ }
+
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 4eb150d..639aded 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
@@ -25,6 +25,8 @@ import it.cnr.isti.workflow.manager.executions.steps.Input;
@Component
public class IteratorContainerExecutor implements ContainerExecutor {
+ private static final long INNER_EXECUTION_WAIT_TIMEOUT_MS = 60_000L;
+
private final ExecutionsService executionsService;
public IteratorContainerExecutor(ExecutionsService executionsService) {
@@ -129,13 +131,11 @@ public class IteratorContainerExecutor implements ContainerExecutor {
private static final Pattern PLACEHOLDER_PATTERN = Pattern.compile("\\$\\{\\{(.*?)}}");
+ private static final long INNER_EXECUTION_WAIT_TIMEOUT_MS = 60_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.
@@ -382,13 +383,11 @@ public class LoopContainerExecutor implements ContainerExecutor {
List findByOwner(String owner);
Page findByOwner(String owner, Pageable pageable);
+ Optional findByIdAndOwner(String id, String owner);
void deleteByOwner(String owner);
}
diff --git a/src/main/resources/application-prod.properties b/src/main/resources/application-prod.properties
new file mode 100644
index 0000000..10501f9
--- /dev/null
+++ b/src/main/resources/application-prod.properties
@@ -0,0 +1,12 @@
+spring.jpa.hibernate.ddl-auto=validate
+app.db.init.enabled=false
+
+# Production security hardening
+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
+
+logging.level.it.cnr.isti.workflow.manager=INFO
+logging.level.root=WARN
diff --git a/src/main/resources/application.properties b/src/main/resources/application.properties
index aa65057..a8328e3 100644
--- a/src/main/resources/application.properties
+++ b/src/main/resources/application.properties
@@ -6,7 +6,7 @@ spring.datasource.username=${DB_USER:lucio}
spring.datasource.password=${DB_PASSWORD:password}
spring.datasource.driver-class-name=org.postgresql.Driver
spring.jpa.properties.hibernate.dialect=org.hibernate.dialect.PostgreSQLDialect
-spring.jpa.hibernate.ddl-auto=${ddl-auto:update}
+spring.jpa.hibernate.ddl-auto=${ddl-auto:validate}
management.endpoints.web.exposure.include=health,info
management.endpoint.health.show-details=always
@@ -28,19 +28,24 @@ management.endpoint.health.show-details=always
# create and drop table, good for testing, production set to none or comment it
-app.db.init.enabled=true
+app.db.init.enabled=${APP_DB_INIT_ENABLED:true}
-app.security.key=${WFEDITOR_SECRET_KEY:088c65fd2a5ca418a79cd10df5dff15c0a79781c0da4fd43c1c14e4e2d7af1ff}
+app.security.key=${WFEDITOR_SECRET_KEY:dev-local-only-change-me-dev-local-only-change-me}
+app.security.default-key=dev-local-only-change-me-dev-local-only-change-me
+app.security.require-explicit-key=${WFEDITOR_REQUIRE_EXPLICIT_SECRET:false}
app.ollama.internal.key=${OLLAMA_INTERNAL_KEY:ollama}
app.ollama.internal.url=${OLLAMA_INTERNAL_URL:https://ollama.internal/api}
app.mcp.bridge.url=${MCP_BRIDGE_URL:http://localhost:8000}
app.mcp.servers.file=${MCP_SERVERS_FILE:}
+app.executions.cache.max-size=${APP_EXECUTIONS_CACHE_MAX_SIZE:1000}
+app.executions.cache.final-ttl-ms=${APP_EXECUTIONS_CACHE_FINAL_TTL_MS:1800000}
+app.executions.cache.cleanup-interval-ms=${APP_EXECUTIONS_CACHE_CLEANUP_INTERVAL_MS:60000}
app.assistant.default-model=${ASSISTANT_DEFAULT_MODEL:gpt-oss:20b}
cors.allowed-origins=${CORS_ALLOWED_ORIGINS:http://localhost:4200}
app.auth.cookie.name=${AUTH_COOKIE_NAME:auth_token}
-app.auth.cookie.secure=${AUTH_COOKIE_SECURE:false}
+app.auth.cookie.secure=${AUTH_COOKIE_SECURE:true}
app.auth.cookie.max-age-seconds=${AUTH_COOKIE_MAX_AGE_SECONDS:86400}
app.import.path=${IMPORT_PATH:/workflow-editor-init}
app.import.enabled=true
diff --git a/src/test/java/it/cnr/isti/workflow/manager/executions/ExecutionsServiceCacheEvictionTest.java b/src/test/java/it/cnr/isti/workflow/manager/executions/ExecutionsServiceCacheEvictionTest.java
new file mode 100644
index 0000000..dbb240d
--- /dev/null
+++ b/src/test/java/it/cnr/isti/workflow/manager/executions/ExecutionsServiceCacheEvictionTest.java
@@ -0,0 +1,106 @@
+package it.cnr.isti.workflow.manager.executions;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+import java.util.Map;
+import java.util.concurrent.ConcurrentMap;
+
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.springframework.test.util.ReflectionTestUtils;
+
+import it.cnr.isti.workflow.manager.executions.repo.ExecutionRepository;
+import it.cnr.isti.workflow.manager.flows.validation.FlowExecutionValidator;
+import it.cnr.isti.workflow.manager.mcp.MCPAgentService;
+
+class ExecutionsServiceCacheEvictionTest {
+
+ private ExecutionsService service;
+
+ @BeforeEach
+ void setup() {
+ service = new ExecutionsService();
+ service.flowExecutionValidator = mock(FlowExecutionValidator.class);
+ service.llmProviders = Map.of();
+ service.executionRepository = mock(ExecutionRepository.class);
+ service.mcpAgentService = mock(MCPAgentService.class);
+
+ ReflectionTestUtils.setField(service, "maxInMemoryExecutions", 2);
+ ReflectionTestUtils.setField(service, "finalStateTtlMs", 100L);
+ }
+
+ @Test
+ void cleanupRemovesExpiredFinalExecution() {
+ ExecutionObject execution = mockExecution("exec-1", ExecutionStatus.SUCCESS);
+ putInCache("exec-1", execution, System.currentTimeMillis() - 500L);
+
+ service.cleanupInMemoryExecutions();
+
+ assertFalse(executions().containsKey("exec-1"));
+ verify(execution).shutdown();
+ }
+
+ @Test
+ void cleanupKeepsExpiredRunningExecution() {
+ ExecutionObject execution = mockExecution("exec-2", ExecutionStatus.RUNNING);
+ putInCache("exec-2", execution, System.currentTimeMillis() - 500L);
+
+ service.cleanupInMemoryExecutions();
+
+ assertTrue(executions().containsKey("exec-2"));
+ }
+
+ @Test
+ void cleanupEvictsOldestFinalWhenCacheExceedsLimit() {
+ ReflectionTestUtils.setField(service, "finalStateTtlMs", Long.MAX_VALUE);
+
+ ExecutionObject oldest = mockExecution("exec-old", ExecutionStatus.SUCCESS);
+ ExecutionObject newest = mockExecution("exec-new", ExecutionStatus.SUCCESS);
+ ExecutionObject running = mockExecution("exec-running", ExecutionStatus.RUNNING);
+
+ long now = System.currentTimeMillis();
+ putInCache("exec-old", oldest, now - 3_000L);
+ putInCache("exec-new", newest, now - 2_000L);
+ putInCache("exec-running", running, now - 1_000L);
+
+ service.cleanupInMemoryExecutions();
+
+ assertEquals(2, executions().size());
+ assertFalse(executions().containsKey("exec-old"));
+ assertTrue(executions().containsKey("exec-new"));
+ assertTrue(executions().containsKey("exec-running"));
+ verify(oldest).shutdown();
+ }
+
+ @SuppressWarnings("unchecked")
+ private ConcurrentMap executions() {
+ return (ConcurrentMap) ReflectionTestUtils.getField(service, "executions");
+ }
+
+ @SuppressWarnings("unchecked")
+ private ConcurrentMap lastAccessMap() {
+ return (ConcurrentMap) ReflectionTestUtils.getField(service, "lastAccessByExecutionId");
+ }
+
+ private void putInCache(String id, ExecutionObject execution, long lastAccessTime) {
+ executions().put(id, execution);
+ lastAccessMap().put(id, lastAccessTime);
+ }
+
+ private ExecutionObject mockExecution(String id, ExecutionStatus status) {
+ ExecutionObject execution = mock(ExecutionObject.class);
+ ExecutionContext context = mock(ExecutionContext.class);
+
+ when(execution.getId()).thenReturn(id);
+ when(execution.getContext()).thenReturn(context);
+ when(context.getStatus()).thenReturn(status);
+ when(context.getExecutionVariableDescriptors()).thenReturn(Map.of());
+
+ return execution;
+ }
+}