diff --git a/src/main/java/it/cnr/isti/workflow/manager/controllers/ApiExceptionHandler.java b/src/main/java/it/cnr/isti/workflow/manager/controllers/ApiExceptionHandler.java index 79913f0..7385490 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/controllers/ApiExceptionHandler.java +++ b/src/main/java/it/cnr/isti/workflow/manager/controllers/ApiExceptionHandler.java @@ -61,7 +61,11 @@ public class ApiExceptionHandler { @ExceptionHandler(ResponseStatusException.class) public ProblemDetail handleResponseStatus(ResponseStatusException e) { - logger.warn("Request failed with status {}: {}", e.getStatusCode(), e.getReason()); + if (e.getStatusCode().is5xxServerError()) { + logger.error("Request failed with status {}: {}", e.getStatusCode(), e.getReason(), e); + } else { + logger.warn("Request failed with status {}: {}", e.getStatusCode(), e.getReason()); + } ProblemDetail problem = ProblemDetail.forStatus(e.getStatusCode()); problem.setTitle(e.getStatusCode().toString()); @@ -72,6 +76,7 @@ public class ApiExceptionHandler { @ExceptionHandler(Exception.class) public ProblemDetail handleGenericException(Exception e) { + logger.error("Unhandled request failure", e); ProblemDetail problem = ProblemDetail.forStatus(HttpStatus.INTERNAL_SERVER_ERROR); problem.setTitle(HttpStatus.INTERNAL_SERVER_ERROR.toString()); problem.setDetail(e.getMessage() == null || e.getMessage().isBlank() ? "Internal server error" : e.getMessage()); diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/MCPAgentChatExecutor.java b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/MCPAgentChatExecutor.java index 8e97c95..d3340c5 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/MCPAgentChatExecutor.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/MCPAgentChatExecutor.java @@ -177,10 +177,13 @@ public class MCPAgentChatExecutor implements BlockExecutor history = existingHistory(partialResults); @@ -197,10 +200,13 @@ public class MCPAgentChatExecutor implements BlockExecutor implements InputListener { listener.completed(this.id, outputs); } catch (Throwable e) { this.status = StepStatus.FAILED; - listener.failed(this.id, e.getMessage()); + String failureMessage = buildFailureMessage(e); + logger.error("Step {} of node {} failed: {}", this.id, this.node.getName(), failureMessage, e); + listener.failed(this.id, failureMessage); } logger.info("Finished executing step " + this.id + " of node " + this.node.getName() + " with status " + this.status); @@ -325,6 +327,30 @@ public class Step implements InputListener { || this.status == StepStatus.CANCELLED; } + private String buildFailureMessage(Throwable error) { + if (error == null) { + return "Unexpected error"; + } + Throwable rootCause = rootCause(error); + String rootMessage = rootCause.getMessage(); + if (rootMessage == null || rootMessage.isBlank()) { + rootMessage = rootCause.getClass().getSimpleName(); + } + String topLevelMessage = error.getMessage(); + if (topLevelMessage == null || topLevelMessage.isBlank() || topLevelMessage.equals(rootMessage)) { + return rootMessage; + } + return topLevelMessage + " | root cause: " + rootMessage; + } + + private Throwable rootCause(Throwable error) { + Throwable current = error; + while (current.getCause() != null && current.getCause() != current) { + current = current.getCause(); + } + return current; + } + public void cancel() { if (isTerminal()) { return; diff --git a/src/main/java/it/cnr/isti/workflow/manager/llms/providers/ollama/InternalOllamaLLMProvider.java b/src/main/java/it/cnr/isti/workflow/manager/llms/providers/ollama/InternalOllamaLLMProvider.java index ec94317..491a48a 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/llms/providers/ollama/InternalOllamaLLMProvider.java +++ b/src/main/java/it/cnr/isti/workflow/manager/llms/providers/ollama/InternalOllamaLLMProvider.java @@ -5,12 +5,15 @@ import java.util.LinkedHashMap; import java.util.List; import java.util.Map; import java.util.Objects; +import java.util.concurrent.TimeoutException; import org.slf4j.Logger; import org.springframework.beans.factory.annotation.Value; import org.springframework.http.MediaType; import org.springframework.stereotype.Service; import org.springframework.web.reactive.function.client.WebClient; +import org.springframework.web.reactive.function.client.WebClientRequestException; +import org.springframework.web.reactive.function.client.WebClientResponseException; import tools.jackson.databind.ObjectMapper; import it.cnr.isti.workflow.manager.llms.ChatMessage; @@ -25,6 +28,7 @@ import reactor.core.publisher.Mono; public class InternalOllamaLLMProvider implements LLMProvider { private static final Logger log = org.slf4j.LoggerFactory.getLogger(InternalOllamaLLMProvider.class); + private static final int MAX_ERROR_BODY_LOG_LENGTH = 1_000; private String ollamaKey; private final WebClient.Builder webClientBuilder; @@ -89,7 +93,7 @@ public class InternalOllamaLLMProvider implements LLMProvider { .bodyToMono(String.class) .timeout(Duration.ofMinutes(2)); - return parseGenerateResponse(result.block(), mapper); + return parseGenerateResponse(blockOllamaCall(result, "generate", model, jsonResponse), mapper); } @Override @@ -123,7 +127,7 @@ public class InternalOllamaLLMProvider implements LLMProvider { .bodyToMono(String.class) .timeout(Duration.ofMinutes(2)); - return parseChatResponse(result.block(), mapper); + return parseChatResponse(blockOllamaCall(result, "chat", model, false), mapper); } private String parseGenerateResponse(String responseBody, ObjectMapper mapper) { @@ -190,7 +194,83 @@ public class InternalOllamaLLMProvider implements LLMProvider { .map(ModelInfo::getName) .toList()); - return result.block(); + return blockOllamaCall(result, "tags", null, false); + } + + private T blockOllamaCall(Mono result, String endpoint, String model, boolean jsonResponse) { + try { + return result.block(); + } catch (RuntimeException e) { + logOllamaFailure(endpoint, model, jsonResponse, e); + throw e; + } + } + + private void logOllamaFailure(String endpoint, String model, boolean jsonResponse, RuntimeException error) { + Throwable root = rootCause(error); + String modelLabel = model == null ? "(none)" : model; + if (error instanceof WebClientResponseException responseError) { + log.error( + "Ollama provider HTTP error. endpoint={}, model={}, jsonResponse={}, status={}, body={}", + endpoint, + modelLabel, + jsonResponse, + responseError.getStatusCode(), + abbreviate(responseError.getResponseBodyAsString()), + error); + return; + } + if (error instanceof WebClientRequestException) { + log.error( + "Ollama provider connection error. endpoint={}, model={}, jsonResponse={}, ollamaUrl={}, errorType={}, message={}", + endpoint, + modelLabel, + jsonResponse, + ollamaURL, + root.getClass().getSimpleName(), + root.getMessage(), + error); + return; + } + if (root instanceof TimeoutException) { + log.error( + "Ollama provider timeout. endpoint={}, model={}, jsonResponse={}, ollamaUrl={}, message={}", + endpoint, + modelLabel, + jsonResponse, + ollamaURL, + root.getMessage(), + error); + return; + } + log.error( + "Ollama provider call failed. endpoint={}, model={}, jsonResponse={}, ollamaUrl={}, errorType={}, message={}", + endpoint, + modelLabel, + jsonResponse, + ollamaURL, + root.getClass().getSimpleName(), + root.getMessage(), + error); + } + + private Throwable rootCause(Throwable error) { + Throwable current = error; + while (current.getCause() != null && current.getCause() != current) { + current = current.getCause(); + } + return current; + } + + private String abbreviate(String value) { + if (value == null || value.isBlank()) { + return "(empty)"; + } + String normalized = value.replaceAll("\\s+", " ").trim(); + if (normalized.length() <= MAX_ERROR_BODY_LOG_LENGTH) { + return normalized; + } + return normalized.substring(0, MAX_ERROR_BODY_LOG_LENGTH) + "..."; } } diff --git a/src/main/java/it/cnr/isti/workflow/manager/mcp/MCPAgentService.java b/src/main/java/it/cnr/isti/workflow/manager/mcp/MCPAgentService.java index 28e03ff..2fb33bf 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/mcp/MCPAgentService.java +++ b/src/main/java/it/cnr/isti/workflow/manager/mcp/MCPAgentService.java @@ -7,6 +7,7 @@ import java.util.LinkedHashMap; import java.util.List; import java.util.Map; import java.util.Objects; +import java.util.concurrent.TimeoutException; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import java.util.concurrent.atomic.AtomicInteger; @@ -15,8 +16,10 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Value; import org.springframework.http.MediaType; +import org.springframework.http.HttpStatusCode; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Component; +import org.springframework.web.reactive.function.client.WebClientRequestException; import org.springframework.web.reactive.function.client.WebClient; import org.springframework.util.StringUtils; @@ -115,8 +118,15 @@ public class MCPAgentService { .retrieve() .onStatus(status -> status.isError(), clientResponse -> clientResponse.bodyToMono(String.class) .defaultIfEmpty("MCP bridge returned an error without body") - .flatMap(body -> Mono.error(new RuntimeException( - "Unable to create MCP bridge session: " + body)))) + .flatMap(body -> { + HttpStatusCode statusCode = clientResponse.statusCode(); + log.warn( + "MCP bridge create session failed: status={} baseUrl={} body={} mcpServersCount={}", + statusCode.value(), mcpBridgeUrl, sanitizeForLog(body), + mcpServers == null ? 0 : mcpServers.size()); + return Mono.error(new RuntimeException( + "Unable to create MCP bridge session: " + body)); + })) .bodyToMono(SessionResponse.class); SessionResponse response = blockWithRetry(requestMono, Duration.ofSeconds(openTimeoutSeconds), "create MCP bridge session"); @@ -161,7 +171,12 @@ public class MCPAgentService { if (server.serverName() != null && !server.serverName().isBlank()) { MCPServersProvider.MCPServerDefinition definition = mcpServersProvider.getServer(server.serverName()); Map serverRequest = new LinkedHashMap<>(); - if (definition.transport() != null && !definition.transport().isBlank()) { + // Do NOT include the "transport" field for stdio servers: the MCP bridge only + // accepts "streamable-http" or "sse" as transport values and rejects "stdio" + // with a validation error. Omitting the field lets the bridge infer the correct + // transport from the presence of "command"/"args", which is the intended behavior. + if (definition.transport() != null && !definition.transport().isBlank() + && !"stdio".equalsIgnoreCase(definition.transport())) { serverRequest.put("transport", definition.transport()); } if (isHttpTransport(definition.transport())) { @@ -212,8 +227,14 @@ public class MCPAgentService { .retrieve() .onStatus(status -> status.isError(), clientResponse -> clientResponse.bodyToMono(String.class) .defaultIfEmpty("MCP bridge returned an error without body") - .flatMap(body -> Mono.error(new RuntimeException( - "Unable to execute MCP bridge query: " + body)))) + .flatMap(body -> { + HttpStatusCode statusCode = clientResponse.statusCode(); + log.warn( + "MCP bridge query failed: status={} sessionId={} baseUrl={} body={}", + statusCode.value(), sessionId, mcpBridgeUrl, sanitizeForLog(body)); + return Mono.error(new RuntimeException( + "Unable to execute MCP bridge query: " + body)); + })) .bodyToMono(QueryResponse.class); QueryResponse response = blockWithRetry(requestMono, Duration.ofSeconds(queryTimeoutSeconds), "execute MCP bridge query"); @@ -235,8 +256,14 @@ public class MCPAgentService { .retrieve() .onStatus(status -> status.isError(), clientResponse -> clientResponse.bodyToMono(String.class) .defaultIfEmpty("MCP bridge returned an error without body") - .flatMap(body -> Mono.error(new RuntimeException( - "Unable to delete MCP bridge session: " + body)))) + .flatMap(body -> { + HttpStatusCode statusCode = clientResponse.statusCode(); + log.warn( + "MCP bridge delete session failed: status={} sessionId={} baseUrl={} body={}", + statusCode.value(), sessionId, mcpBridgeUrl, sanitizeForLog(body)); + return Mono.error(new RuntimeException( + "Unable to delete MCP bridge session: " + body)); + })) .toBodilessEntity() .then(); blockWithRetry(requestMono, Duration.ofSeconds(closeTimeoutSeconds), "delete MCP bridge session"); @@ -273,6 +300,16 @@ public class MCPAgentService { onSuccess(); return response; } catch (Exception ex) { + Throwable rootCause = rootCause(ex); + log.error( + "MCP bridge operation failed: operation={} baseUrl={} timeoutSeconds={} retryAttempts={} failureType={} message={}", + operationLabel, + mcpBridgeUrl, + timeout.toSeconds(), + Math.max(0, retryAttempts), + classifyFailure(rootCause), + sanitizeForLog(rootCause.getMessage()), + ex); onFailure(ex); throw new RuntimeException("Unable to " + operationLabel, ex); } @@ -303,6 +340,49 @@ public class MCPAgentService { } } + private Throwable rootCause(Throwable throwable) { + Throwable current = throwable; + while (current.getCause() != null && current.getCause() != current) { + current = current.getCause(); + } + return current; + } + + private String classifyFailure(Throwable throwable) { + if (throwable == null) { + return "unknown"; + } + if (throwable instanceof TimeoutException) { + return "timeout"; + } + if (throwable instanceof WebClientRequestException) { + return "request-connection"; + } + String message = throwable.getMessage(); + if (message != null) { + String lower = message.toLowerCase(); + if (lower.contains(" 4") || lower.contains("status=4")) { + return "http-4xx"; + } + if (lower.contains(" 5") || lower.contains("status=5")) { + return "http-5xx"; + } + } + return throwable.getClass().getSimpleName(); + } + + private String sanitizeForLog(String value) { + if (value == null) { + return ""; + } + String normalized = value.replaceAll("[\\r\\n\\t]+", " ").trim(); + int maxLen = 1000; + if (normalized.length() <= maxLen) { + return normalized; + } + return normalized.substring(0, maxLen) + "..."; + } + private boolean isHttpTransport(String transport) { return "http".equalsIgnoreCase(transport) || "https".equalsIgnoreCase(transport)