From 48925e19d4cff94e35dbfea612ffb7178392d62d Mon Sep 17 00:00:00 2001 From: Lucio Lelii Date: Tue, 5 May 2026 19:53:23 +0200 Subject: [PATCH] Harden config and improve execution runtime resilience --- pom.xml | 28 +++- .../workflow/manager/auth/config/JwtUtil.java | 7 +- .../controllers/ExecutionsController.java | 6 +- .../manager/executions/ExecutionContext.java | 28 ++++ .../manager/executions/ExecutionObject.java | 15 +- .../manager/executions/ExecutionsService.java | 130 +++++++++++++++++- .../containers/IteratorContainerExecutor.java | 12 +- .../containers/LoopContainerExecutor.java | 11 +- .../executions/repo/ExecutionRepository.java | 2 + .../resources/application-prod.properties | 12 ++ src/main/resources/application.properties | 13 +- .../ExecutionsServiceCacheEvictionTest.java | 106 ++++++++++++++ 12 files changed, 342 insertions(+), 28 deletions(-) create mode 100644 src/main/resources/application-prod.properties create mode 100644 src/test/java/it/cnr/isti/workflow/manager/executions/ExecutionsServiceCacheEvictionTest.java 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; + } +}