From e1788a86c2bcaa7401f6cf7c0081a0f93bd96c96 Mon Sep 17 00:00:00 2001 From: Lucio Lelii Date: Mon, 14 Apr 2025 17:03:44 +0200 Subject: [PATCH] updated --- .../workflow/manager/ExecutionService.java | 9 ++++-- .../manager/executors/ExecutionObject.java | 2 -- .../ai/services/OpenRouterService.java | 8 +++++- .../manager/model/ExecutionContext.java | 28 +++++++++++++++---- 4 files changed, 36 insertions(+), 11 deletions(-) diff --git a/src/main/java/it/cnr/isti/workflow/manager/ExecutionService.java b/src/main/java/it/cnr/isti/workflow/manager/ExecutionService.java index 4d13a3d..f454c0c 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/ExecutionService.java +++ b/src/main/java/it/cnr/isti/workflow/manager/ExecutionService.java @@ -23,7 +23,6 @@ public class ExecutionService { public ExecutionObject createExecution(FlowEntity flow) { ExecutionObject execObject = transformerService.transform(flow.getFlow()); - executions.put(execObject.getId(), execObject); return execObject; } @@ -43,8 +42,14 @@ public class ExecutionService { public ExecutionObject prepareInput(String id, String key, Object input) { ExecutionObject eo = getExecution(id); if (eo == null) throw new IllegalArgumentException("Execution with id "+id+" not found"); + + if (eo.getContext().getStatus().isInitState()) { + throw new IllegalStateException("Execution with id "+id+" is not in initialization status (CURRENT STATUS is "+eo.getContext().getStatus()+")"); + } + eo.getContext().setStatus(Status.INITIALIZING); eo.getStartStep().inputReady(key, input); + eo.getContext().addInput(key, input); if (eo.getStartStep().areAllInputsReady()) { eo.getContext().setStatus(Status.READY); } @@ -55,7 +60,7 @@ public class ExecutionService { ExecutionObject eo = getExecution(id); if (eo == null) throw new IllegalArgumentException("Execution with id "+id+" not found"); if (eo.getContext().getStatus() != Status.READY) { - throw new IllegalStateException("Execution with id "+id+" is not ready"); + throw new IllegalStateException("Execution with id "+id+" is not in READY status (CURRENT STATUS is "+eo.getContext().getStatus()+")"); } eo.getStartStep().startExecution(); eo.getContext().setStatus(Status.RUNNING); diff --git a/src/main/java/it/cnr/isti/workflow/manager/executors/ExecutionObject.java b/src/main/java/it/cnr/isti/workflow/manager/executors/ExecutionObject.java index 50e6286..7724f4e 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executors/ExecutionObject.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executors/ExecutionObject.java @@ -40,8 +40,6 @@ public class ExecutionObject implements InputReadyListener { @JsonIgnore ExecutionStep endStep; - @JsonIgnore - Map inputs = new HashMap<>(); public ExecutionObject(Flow flow) { this.flow = flow; diff --git a/src/main/java/it/cnr/isti/workflow/manager/executors/ai/services/OpenRouterService.java b/src/main/java/it/cnr/isti/workflow/manager/executors/ai/services/OpenRouterService.java index 470514a..4556907 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executors/ai/services/OpenRouterService.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executors/ai/services/OpenRouterService.java @@ -2,10 +2,14 @@ package it.cnr.isti.workflow.manager.executors.ai.services; import org.springframework.web.reactive.function.client.WebClient; import org.springframework.http.MediaType; +import org.springframework.http.client.reactive.ReactorClientHttpConnector; import org.springframework.stereotype.Service; import reactor.core.publisher.Mono; +import reactor.netty.http.client.HttpClient; + import java.util.Map; +import java.time.Duration; import java.util.List; import java.util.Objects; @@ -17,7 +21,9 @@ public class OpenRouterService { private static final String OpenRouter_URL = "https://openrouter.ai/api/v1/chat/completions"; public OpenRouterService(WebClient.Builder webClientBuilder) { - this.webClientBuilder = webClientBuilder; + HttpClient client = HttpClient.create() + .responseTimeout(Duration.ofMinutes(2)); + this.webClientBuilder = webClientBuilder.clientConnector(new ReactorClientHttpConnector(client)); } public void setApiKey(String apiKey) { diff --git a/src/main/java/it/cnr/isti/workflow/manager/model/ExecutionContext.java b/src/main/java/it/cnr/isti/workflow/manager/model/ExecutionContext.java index 48bbe95..3da4afa 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/model/ExecutionContext.java +++ b/src/main/java/it/cnr/isti/workflow/manager/model/ExecutionContext.java @@ -14,14 +14,26 @@ import lombok.Setter; public class ExecutionContext { public enum Status { - CREATED, - INITIALIZING, - READY, - RUNNING, - SUCCESS, - ERROR + CREATED(true), + INITIALIZING(true), + READY(true), + RUNNING(true), + SUCCESS(true), + ERROR(true); + + private boolean initState; + + Status(boolean initState) { + this.initState = initState; + } + + public boolean isInitState() { + return initState; + } } + Map inputs = new HashMap<>(); + @Setter Map result = null; @@ -44,6 +56,10 @@ public class ExecutionContext { this.endTime = System.currentTimeMillis(); } } + + public void addInput(String key, Object value) { + this.inputs.put(key, value); + } }