package com.wok.supportbot.config;
import com.wok.supportbot.entity.McpServerConfig;
import com.wok.supportbot.service.McpServerConfigService;
import io.modelcontextprotocol.client.McpSyncClient;
import io.modelcontextprotocol.client.transport.HttpClientSseClientTransport;
import io.modelcontextprotocol.client.transport.ServerParameters;
import io.modelcontextprotocol.client.transport.StdioClientTransport;
import io.modelcontextprotocol.spec.McpSchema;
import jakarta.annotation.PostConstruct;
import jakarta.annotation.PreDestroy;
import lombok.extern.slf4j.Slf4j;
import org.springframework.ai.model.ModelOptionsUtils;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component;
import java.time.Duration;
import java.time.Instant;
import java.util.*;
import java.util.concurrent.ConcurrentHashMap;
/**
* MCP 客户端生命周期管理器
* 负责 MCP Client 的创建、缓存、刷新和销毁。
*
* 支持两种传输模式:
* - SSE(Server-Sent Events):通过 HttpClientSseClientTransport 连接远程 MCP Server
* - stdio(标准输入输出):通过 StdioClientTransport 启动本地 MCP Server 进程
*/
@Component
@Slf4j
public class McpClientManager {
@Autowired
private McpServerConfigService mcpServerConfigService;
/**
* 客户端缓存:key = 配置ID 的字符串形式,value = MCP Client 实例
* volatile 保证 refreshAll() 双缓冲切换时的可见性
*/
private volatile ConcurrentHashMap clientCache = new ConcurrentHashMap<>();
/**
* 不可用配置集合:记录已知不存在或未启用的配置ID,避免重复查询 DB
* 配置被启用或新建时需同步移除此集合中的对应条目
*/
private final Set unavailableConfigs = ConcurrentHashMap.newKeySet();
/**
* 健康状态缓存:key = 配置ID 的字符串形式,value = 健康检查结果
*/
private final ConcurrentHashMap healthCache = new ConcurrentHashMap<>();
/**
* MCP Server 健康状态记录
*/
public record HealthStatus(
/** 状态:ONLINE / OFFLINE / UNKNOWN */
String status,
/** 响应延迟(毫秒) */
long latencyMs,
/** 最后检查时间 */
Instant lastCheckTime,
/** 错误信息(健康时为 null) */
String errorMessage
) {
public static HealthStatus online(long latencyMs) {
return new HealthStatus("ONLINE", latencyMs, Instant.now(), null);
}
public static HealthStatus offline(long latencyMs, String errorMessage) {
return new HealthStatus("OFFLINE", latencyMs, Instant.now(), errorMessage);
}
public static HealthStatus unknown(String errorMessage) {
return new HealthStatus("UNKNOWN", 0, Instant.now(), errorMessage);
}
}
/**
* MCP 客户端请求超时时间,支持 Duration 格式(如 30s、PT30S、60s)
*/
@Value("${mcp.client.request-timeout:30s}")
private String requestTimeoutStr;
/**
* MCP 客户端初始化超时时间,支持 Duration 格式(如 15s、PT15S、30s)
*/
@Value("${mcp.client.init-timeout:15s}")
private String initTimeoutStr;
/**
* 应用启动时自动加载所有已启用的 MCP Server 配置并建立连接
* 确保第一次对话请求时 clientCache 已就绪,MCP 工具可以被注册到 ChatClient
*/
@PostConstruct
public void init() {
log.info("McpClientManager 初始化,开始加载已启用的 MCP Server 配置...");
refreshAll();
}
/**
* 刷新所有 MCP 客户端连接(双缓冲策略)
* 先在新 Map 中构建所有客户端,再原子切换引用,最后关闭旧客户端。
* 避免清空与重建之间的请求全部失败。
*/
public void refreshAll() {
log.info("开始刷新所有 MCP 客户端连接...");
// 获取所有启用的配置(一次查询代替两次按类型过滤)
List allActiveConfigs = mcpServerConfigService.listAllActiveConfigs();
// 第一步:在全新 Map 中构建所有客户端(不影响现有缓存)
ConcurrentHashMap newCache = new ConcurrentHashMap<>();
for (McpServerConfig config : allActiveConfigs) {
try {
McpSyncClient client = createClientDirectly(config);
if (client != null) {
newCache.put(config.getId().toString(), client);
}
} catch (Exception e) {
log.error("创建客户端失败: id={}, name={}, transportType={}, error={}",
config.getId(), config.getName(), config.getTransportType(), e.getMessage());
}
}
// 第二步:原子切换引用(volatile 保证其他线程立即可见)
ConcurrentHashMap oldCache = clientCache;
clientCache = newCache;
unavailableConfigs.clear();
// 第三步:关闭旧缓存中的客户端(不影响新请求)
for (Map.Entry entry : oldCache.entrySet()) {
try {
McpSyncClient client = entry.getValue();
if (client != null) {
client.close();
log.debug("已关闭旧 MCP 客户端: configId={}", entry.getKey());
}
} catch (Exception e) {
log.error("关闭旧 MCP 客户端失败: configId={}, error={}", entry.getKey(), e.getMessage());
}
}
log.info("MCP 客户端刷新完成,当前缓存数量: {}", clientCache.size());
}
/**
* 获取指定配置的客户端(懒加载)
* 使用 computeIfAbsent 保证同一配置只创建一次客户端,避免并发竞态。
* 如果缓存中不存在且不在不可用集合中,从 DB 读取配置后原子创建并缓存。
*
* @param configId 配置ID
* @return MCP Client 实例,不存在或未启用返回 null
*/
public McpSyncClient getClient(Long configId) {
String key = configId.toString();
// 快速路径:缓存命中
McpSyncClient cached = clientCache.get(key);
if (cached != null) {
return cached;
}
// 快速路径:已知不可用的配置,直接跳过
if (unavailableConfigs.contains(key)) {
return null;
}
// 缓存未命中,使用 computeIfAbsent 保证同一 key 只创建一次客户端
// ConcurrentHashMap.computeIfAbsent 对同一 key 加锁,避免并发线程重复创建
// 注意:mapping function 不能返回 null(会抛 NPE),因此不可用的情况记录到 unavailableConfigs
McpServerConfig config = mcpServerConfigService.getConfigById(configId);
if (config == null) {
log.warn("MCP 配置不存在: id={}", configId);
unavailableConfigs.add(key);
return null;
}
if (!Boolean.TRUE.equals(config.getIsActive())) {
log.warn("MCP 配置未启用,跳过创建客户端: id={}, name={}", configId, config.getName());
unavailableConfigs.add(key);
return null;
}
// 配置存在且启用,通过 computeIfAbsent 原子创建(防止并发重复创建)
return clientCache.computeIfAbsent(key, k -> createClientDirectly(config));
}
/**
* 获取所有已启用的 MCP 工具描述
* 遍历所有已缓存的客户端,调用 listTools() 获取真实的工具列表
*
* @return 工具描述列表,每个元素包含 config_id、name、transport_type、description、tools
*/
public List