Harden config and improve execution runtime resilience
This commit is contained in:
parent
2c5df42326
commit
48925e19d4
28
pom.xml
28
pom.xml
|
|
@ -97,6 +97,10 @@
|
|||
<groupId>org.postgresql</groupId>
|
||||
<artifactId>postgresql</artifactId>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>org.flywaydb</groupId>
|
||||
<artifactId>flyway-core</artifactId>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>org.springframework.boot</groupId>
|
||||
<artifactId>spring-boot-starter-web</artifactId>
|
||||
|
|
@ -168,6 +172,10 @@
|
|||
<groupId>org.apache.maven.plugins</groupId>
|
||||
<artifactId>maven-compiler-plugin</artifactId>
|
||||
<configuration>
|
||||
<release>${java.version}</release>
|
||||
<compilerArgs>
|
||||
<arg>--enable-preview</arg>
|
||||
</compilerArgs>
|
||||
<annotationProcessorPaths>
|
||||
<path>
|
||||
<groupId>org.projectlombok</groupId>
|
||||
|
|
@ -192,12 +200,30 @@
|
|||
<groupId>org.apache.maven.plugins</groupId>
|
||||
<artifactId>maven-surefire-plugin</artifactId>
|
||||
<configuration>
|
||||
<argLine>-javaagent:${settings.localRepository}/net/bytebuddy/byte-buddy-agent/${byte-buddy.version}/byte-buddy-agent-${byte-buddy.version}.jar</argLine>
|
||||
<argLine>--enable-preview ${argLine} -javaagent:${settings.localRepository}/net/bytebuddy/byte-buddy-agent/${byte-buddy.version}/byte-buddy-agent-${byte-buddy.version}.jar</argLine>
|
||||
<systemPropertyVariables>
|
||||
<mockito.mock-maker>subclass</mockito.mock-maker>
|
||||
</systemPropertyVariables>
|
||||
</configuration>
|
||||
</plugin>
|
||||
<plugin>
|
||||
<groupId>org.jacoco</groupId>
|
||||
<artifactId>jacoco-maven-plugin</artifactId>
|
||||
<executions>
|
||||
<execution>
|
||||
<goals>
|
||||
<goal>prepare-agent</goal>
|
||||
</goals>
|
||||
</execution>
|
||||
<execution>
|
||||
<id>report</id>
|
||||
<phase>test</phase>
|
||||
<goals>
|
||||
<goal>report</goal>
|
||||
</goals>
|
||||
</execution>
|
||||
</executions>
|
||||
</plugin>
|
||||
|
||||
</plugins>
|
||||
</build>
|
||||
|
|
|
|||
|
|
@ -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());
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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());
|
||||
}
|
||||
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<String, Object> runtimeContextValues) {
|
||||
|
|
|
|||
|
|
@ -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<Step<?>> 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<String, Object> runtimeContext = ExecutionRuntimeContextSupport.contextValues(
|
||||
this.id,
|
||||
|
|
|
|||
|
|
@ -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<String, ExecutionObject> executions = new ConcurrentHashMap<>();
|
||||
private final ConcurrentMap<String, Long> 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<ExecutionAuthorizationRequirement> 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<ExecutionObject> getAllExecutions() {
|
||||
return executionRepository.findAll().stream()
|
||||
.map(entity -> executions.computeIfAbsent(entity.getId(), ignored -> rebuildExecution(entity)))
|
||||
List<ExecutionObject> 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<ExecutionObject> getExecutionsByOwner(String owner, Pageable pageable) {
|
||||
return executionRepository.findByOwner(owner, pageable)
|
||||
.map(entity -> executions.computeIfAbsent(entity.getId(), ignored -> rebuildExecution(entity)));
|
||||
Page<ExecutionObject> 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<ExecutionAuthorizationRequirement> 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<String> 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<String> 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;
|
||||
|
|
|
|||
|
|
@ -25,6 +25,8 @@ import it.cnr.isti.workflow.manager.executions.steps.Input;
|
|||
@Component
|
||||
public class IteratorContainerExecutor implements ContainerExecutor<IteratorContainerType> {
|
||||
|
||||
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<IteratorCont
|
|||
private ExecutionObject startAndWait(ExecutionObject innerExecution) {
|
||||
innerExecution = executionsService.startExecution(innerExecution.getId());
|
||||
while (innerExecution.getContext().getStatus() == ExecutionStatus.RUNNING) {
|
||||
try {
|
||||
Thread.sleep(25L);
|
||||
} catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
throw new IllegalStateException("Interrupted while executing IteratorContainer subflow", e);
|
||||
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");
|
||||
}
|
||||
innerExecution = executionsService.getExecution(innerExecution.getId());
|
||||
}
|
||||
|
||||
if (innerExecution.getContext().getStatus() == ExecutionStatus.WAITING) {
|
||||
|
|
|
|||
|
|
@ -40,6 +40,7 @@ import it.cnr.isti.workflow.manager.llms.providers.LLMProvider;
|
|||
public class LoopContainerExecutor implements ContainerExecutor<LoopContainerType> {
|
||||
|
||||
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<LoopContainerTyp
|
|||
private ExecutionObject startAndWait(ExecutionObject innerExecution) {
|
||||
innerExecution = executionsService.startExecution(innerExecution.getId());
|
||||
while (innerExecution.getContext().getStatus() == ExecutionStatus.RUNNING) {
|
||||
try {
|
||||
Thread.sleep(25L);
|
||||
} catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
throw new IllegalStateException("Interrupted while executing LoopContainer subflow", e);
|
||||
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");
|
||||
}
|
||||
innerExecution = executionsService.getExecution(innerExecution.getId());
|
||||
}
|
||||
|
||||
if (innerExecution.getContext().getStatus() == ExecutionStatus.WAITING) {
|
||||
|
|
|
|||
|
|
@ -5,9 +5,11 @@ import org.springframework.data.domain.Pageable;
|
|||
import org.springframework.data.jpa.repository.JpaRepository;
|
||||
|
||||
import java.util.List;
|
||||
import java.util.Optional;
|
||||
|
||||
public interface ExecutionRepository extends JpaRepository<ExecutionEntity, String> {
|
||||
List<ExecutionEntity> findByOwner(String owner);
|
||||
Page<ExecutionEntity> findByOwner(String owner, Pageable pageable);
|
||||
Optional<ExecutionEntity> findByIdAndOwner(String id, String owner);
|
||||
void deleteByOwner(String owner);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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<String, ExecutionObject> executions() {
|
||||
return (ConcurrentMap<String, ExecutionObject>) ReflectionTestUtils.getField(service, "executions");
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
private ConcurrentMap<String, Long> lastAccessMap() {
|
||||
return (ConcurrentMap<String, Long>) 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;
|
||||
}
|
||||
}
|
||||
Loading…
Reference in New Issue