Improve provider failure diagnostics

This commit is contained in:
Lucio Lelii 2026-05-19 13:39:17 +02:00
parent d466459a27
commit 4bd240ddbb
5 changed files with 215 additions and 18 deletions

View File

@ -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());

View File

@ -177,10 +177,13 @@ public class MCPAgentChatExecutor implements BlockExecutor<MCPAgentChatBlockType
+ MCPAgentChatBlockFactory.FINAL_RESPONSE_FIELD);
}
if (!isManagedSharedSession(configuration, executionVariableDescriptors)) {
mcpAgentService.closeSessionQuietly(existingSessionId(partialResults));
if (eventLogger != null) {
String sessionId = existingSessionId(partialResults);
if (StringUtils.hasText(sessionId)) {
mcpAgentService.closeSessionQuietly(sessionId);
}
if (eventLogger != null && StringUtils.hasText(sessionId)) {
eventLogger.info(ExecutionEventType.MCP_SESSION_CLOSED, "Closed MCP session",
Map.of("sessionId", existingSessionId(partialResults)));
Map.of("sessionId", sessionId));
}
}
List<String> history = existingHistory(partialResults);
@ -197,10 +200,13 @@ public class MCPAgentChatExecutor implements BlockExecutor<MCPAgentChatBlockType
MCPAgentChatBlockConfiguration configuration =
(MCPAgentChatBlockConfiguration) block.getSpecificConfiguration();
if (!isManagedSharedSession(configuration, executionVariableDescriptors)) {
mcpAgentService.closeSessionQuietly(existingSessionId(partialResults));
if (eventLogger != null) {
String sessionId = existingSessionId(partialResults);
if (StringUtils.hasText(sessionId)) {
mcpAgentService.closeSessionQuietly(sessionId);
}
if (eventLogger != null && StringUtils.hasText(sessionId)) {
eventLogger.info(ExecutionEventType.MCP_SESSION_CLOSED, "Closed MCP session on cancel",
Map.of("sessionId", existingSessionId(partialResults)));
Map.of("sessionId", sessionId));
}
}
}

View File

@ -182,7 +182,9 @@ public class Step<N extends FlowNode> 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<N extends FlowNode> 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;

View File

@ -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> T blockOllamaCall(Mono<T> 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) + "...";
}
}

View File

@ -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<String, Object> 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)