fix: Deep Agent Builder 호환성을 위한 mcp-http Warm Pool 세션 버그 및 CORS 헤더 수정
This commit is contained in:
@@ -27,6 +27,7 @@ public class CorsConfig implements WebMvcConfigurer {
|
|||||||
.allowedOriginPatterns("*") // 외부 Agent Builder 등 모든 오리진 허용
|
.allowedOriginPatterns("*") // 외부 Agent Builder 등 모든 오리진 허용
|
||||||
.allowedMethods("GET", "POST", "PUT", "DELETE", "OPTIONS", "HEAD", "PATCH") // 허용할 HTTP 메서드
|
.allowedMethods("GET", "POST", "PUT", "DELETE", "OPTIONS", "HEAD", "PATCH") // 허용할 HTTP 메서드
|
||||||
.allowedHeaders("*") // 모든 헤더 허용
|
.allowedHeaders("*") // 모든 헤더 허용
|
||||||
|
.exposedHeaders("Mcp-Session-Id") // MCP-HTTP 세션 아이디 노출 허용
|
||||||
.allowCredentials(true) // 쿠키/인증 정보 허용
|
.allowCredentials(true) // 쿠키/인증 정보 허용
|
||||||
.maxAge(3600); // preflight 요청 캐시 시간 (초 단위)
|
.maxAge(3600); // preflight 요청 캐시 시간 (초 단위)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -48,6 +48,7 @@ public class WebConfig implements WebMvcConfigurer {
|
|||||||
.allowedOriginPatterns("*")
|
.allowedOriginPatterns("*")
|
||||||
.allowedMethods("GET", "POST", "PUT", "DELETE", "OPTIONS")
|
.allowedMethods("GET", "POST", "PUT", "DELETE", "OPTIONS")
|
||||||
.allowedHeaders("*")
|
.allowedHeaders("*")
|
||||||
|
.exposedHeaders("Mcp-Session-Id")
|
||||||
.allowCredentials(true);
|
.allowCredentials(true);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -49,6 +49,7 @@ public class CustomWebMvcSseServerTransportProvider implements McpServerTranspor
|
|||||||
private final String sseEndpoint;
|
private final String sseEndpoint;
|
||||||
private final String messageEndpoint;
|
private final String messageEndpoint;
|
||||||
private final Map<String, McpServerSession> sessions = new ConcurrentHashMap<>();
|
private final Map<String, McpServerSession> sessions = new ConcurrentHashMap<>();
|
||||||
|
private final Map<String, CustomMcpSessionTransport> customTransports = new ConcurrentHashMap<>();
|
||||||
private final ObjectMapper objectMapper;
|
private final ObjectMapper objectMapper;
|
||||||
|
|
||||||
public CustomWebMvcSseServerTransportProvider(String sseEndpoint, String messageEndpoint, ObjectMapper objectMapper) {
|
public CustomWebMvcSseServerTransportProvider(String sseEndpoint, String messageEndpoint, ObjectMapper objectMapper) {
|
||||||
@@ -71,6 +72,7 @@ public class CustomWebMvcSseServerTransportProvider implements McpServerTranspor
|
|||||||
} catch (Exception ignored) {}
|
} catch (Exception ignored) {}
|
||||||
});
|
});
|
||||||
sessions.clear();
|
sessions.clear();
|
||||||
|
customTransports.clear();
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -108,27 +110,38 @@ public class CustomWebMvcSseServerTransportProvider implements McpServerTranspor
|
|||||||
return emitter;
|
return emitter;
|
||||||
}
|
}
|
||||||
|
|
||||||
public SseEmitter handleCustomSse(String sessionId, String body) {
|
public org.springframework.web.servlet.mvc.method.annotation.SseEmitter handleCustomSse(String sessionId, String body) {
|
||||||
if (sessionFactory == null) {
|
if (sessionFactory == null) {
|
||||||
throw new IllegalStateException("SessionFactory not configured");
|
throw new IllegalStateException("SessionFactory not configured");
|
||||||
}
|
}
|
||||||
|
|
||||||
SseEmitter emitter = new SseEmitter(-1L);
|
org.springframework.web.servlet.mvc.method.annotation.SseEmitter emitter = new org.springframework.web.servlet.mvc.method.annotation.SseEmitter(-1L);
|
||||||
|
|
||||||
CustomMcpSessionTransport sessionTransport = new CustomMcpSessionTransport(emitter, sessionId);
|
boolean isNew = !sessions.containsKey(sessionId);
|
||||||
McpServerSession session = sessionFactory.create(sessionTransport);
|
|
||||||
sessions.put(sessionId, session);
|
|
||||||
|
|
||||||
emitter.onCompletion(() -> sessions.remove(sessionId));
|
if (isNew) {
|
||||||
emitter.onTimeout(() -> sessions.remove(sessionId));
|
CustomMcpSessionTransport sessionTransport = new CustomMcpSessionTransport(emitter, sessionId);
|
||||||
|
customTransports.put(sessionId, sessionTransport);
|
||||||
|
McpServerSession session = sessionFactory.create(sessionTransport);
|
||||||
|
sessions.put(sessionId, session);
|
||||||
|
|
||||||
|
// 주의: 클라이언트가 단일 POST 응답 후 연결을 끊더라도,
|
||||||
|
// 웜 풀(Warm Pool) 스펙상 세션은 살려둬야 하므로 세션 삭제 로직 제외
|
||||||
|
} else {
|
||||||
|
// 기존 세션인 경우 Emitter 파이프만 덮어씌움 (Switching)
|
||||||
|
CustomMcpSessionTransport sessionTransport = customTransports.get(sessionId);
|
||||||
|
if (sessionTransport != null) {
|
||||||
|
sessionTransport.setEmitter(emitter);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
new Thread(() -> {
|
new Thread(() -> {
|
||||||
try {
|
try {
|
||||||
// 커스텀 클라이언트는 endpoint 이벤트를 무시할 수 있지만, 표준 호환성을 위해 전송
|
// 커스텀 클라이언트는 endpoint 이벤트를 무시할 수 있지만 표준 호환성을 위해 전송
|
||||||
Thread.sleep(100);
|
Thread.sleep(100);
|
||||||
emitter.send(SseEmitter.event().name("endpoint").data(messageEndpoint + "?sessionId=" + sessionId));
|
emitter.send(org.springframework.web.servlet.mvc.method.annotation.SseEmitter.event().name("endpoint").data(messageEndpoint + "?sessionId=" + sessionId));
|
||||||
|
|
||||||
// Body로 들어온 initialize 등 즉시 처리
|
// Body로 들어온 메시지 즉시 비동기 처리
|
||||||
if (body != null && !body.trim().isEmpty()) {
|
if (body != null && !body.trim().isEmpty()) {
|
||||||
handleMessage(sessionId, body);
|
handleMessage(sessionId, body);
|
||||||
}
|
}
|
||||||
@@ -140,10 +153,10 @@ public class CustomWebMvcSseServerTransportProvider implements McpServerTranspor
|
|||||||
return emitter;
|
return emitter;
|
||||||
}
|
}
|
||||||
|
|
||||||
public ResponseEntity<String> handleMessage(String sessionId, String body) {
|
public org.springframework.http.ResponseEntity<String> handleMessage(String sessionId, String body) {
|
||||||
log.info("Received POST message for sessionId: " + sessionId + ", body: " + body);
|
log.info("Received POST message for sessionId: " + sessionId + ", body: " + body);
|
||||||
if (sessionId == null || !sessions.containsKey(sessionId)) {
|
if (sessionId == null || !sessions.containsKey(sessionId)) {
|
||||||
return ResponseEntity.badRequest().body("Missing or invalid sessionId");
|
return ResponseEntity.badRequest().body("Unexpected request body");
|
||||||
}
|
}
|
||||||
|
|
||||||
McpServerSession session = sessions.get(sessionId);
|
McpServerSession session = sessions.get(sessionId);
|
||||||
@@ -171,15 +184,23 @@ public class CustomWebMvcSseServerTransportProvider implements McpServerTranspor
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public boolean hasSession(String sessionId) {
|
||||||
|
return sessionId != null && sessions.containsKey(sessionId);
|
||||||
|
}
|
||||||
|
|
||||||
private class CustomMcpSessionTransport implements McpServerTransport {
|
private class CustomMcpSessionTransport implements McpServerTransport {
|
||||||
private final SseEmitter emitter;
|
private volatile org.springframework.web.servlet.mvc.method.annotation.SseEmitter emitter;
|
||||||
private final String sessionId;
|
private final String sessionId;
|
||||||
|
|
||||||
public CustomMcpSessionTransport(SseEmitter emitter, String sessionId) {
|
public CustomMcpSessionTransport(org.springframework.web.servlet.mvc.method.annotation.SseEmitter emitter, String sessionId) {
|
||||||
this.emitter = emitter;
|
this.emitter = emitter;
|
||||||
this.sessionId = sessionId;
|
this.sessionId = sessionId;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public void setEmitter(org.springframework.web.servlet.mvc.method.annotation.SseEmitter emitter) {
|
||||||
|
this.emitter = emitter;
|
||||||
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public Mono<Void> sendMessage(McpSchema.JSONRPCMessage message) {
|
public Mono<Void> sendMessage(McpSchema.JSONRPCMessage message) {
|
||||||
return Mono.fromRunnable(() -> {
|
return Mono.fromRunnable(() -> {
|
||||||
@@ -187,18 +208,30 @@ public class CustomWebMvcSseServerTransportProvider implements McpServerTranspor
|
|||||||
try {
|
try {
|
||||||
String json = objectMapper.writeValueAsString(message);
|
String json = objectMapper.writeValueAsString(message);
|
||||||
log.info("Serialized message: " + json);
|
log.info("Serialized message: " + json);
|
||||||
emitter.send(SseEmitter.event().name("message").data(json));
|
if (this.emitter != null) {
|
||||||
log.info("Message successfully sent to SSE emitter");
|
this.emitter.send(org.springframework.web.servlet.mvc.method.annotation.SseEmitter.event().name("message").data(json));
|
||||||
|
log.info("Message successfully sent to SSE emitter");
|
||||||
|
|
||||||
|
// Custom 프로토콜: 1회 요청당 1응답 후 종료 (스트림을 닫아버림)
|
||||||
|
// 클라이언트가 한 번의 POST 후 응답을 받고 연결을 끊기 때문
|
||||||
|
this.emitter.complete();
|
||||||
|
}
|
||||||
} catch (Exception e) {
|
} catch (Exception e) {
|
||||||
log.error("Error sending message to SSE emitter", e);
|
log.error("Error sending message to SSE emitter", e);
|
||||||
emitter.completeWithError(e);
|
if (this.emitter != null) {
|
||||||
|
this.emitter.completeWithError(e);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public Mono<Void> closeGracefully() {
|
public Mono<Void> closeGracefully() {
|
||||||
return Mono.fromRunnable(emitter::complete);
|
return Mono.fromRunnable(() -> {
|
||||||
|
if (this.emitter != null) {
|
||||||
|
this.emitter.complete();
|
||||||
|
}
|
||||||
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
|
|||||||
@@ -70,18 +70,18 @@ public class DynamicMcpController {
|
|||||||
return ResponseEntity.badRequest().body("Unknown category: " + category);
|
return ResponseEntity.badRequest().body("Unknown category: " + category);
|
||||||
}
|
}
|
||||||
|
|
||||||
if (sessionId == null || sessionId.isEmpty()) {
|
boolean isNew = (sessionId == null || sessionId.isEmpty());
|
||||||
// 새 세션 생성 (initialize 요청)
|
String activeSessionId = isNew ? java.util.UUID.randomUUID().toString() : sessionId;
|
||||||
String newSessionId = UUID.randomUUID().toString();
|
|
||||||
SseEmitter emitter = transport.handleCustomSse(newSessionId, body);
|
|
||||||
|
|
||||||
return ResponseEntity.ok()
|
if (!isNew && !transport.hasSession(activeSessionId)) {
|
||||||
.header("Mcp-Session-Id", newSessionId)
|
// 클라이언트가 보낸 세션 ID가 만료되었거나 존재하지 않는 경우 (Warm Pool 스펙: 404 Not Found 반환)
|
||||||
.body(emitter);
|
return ResponseEntity.notFound().build();
|
||||||
} else {
|
|
||||||
// 기존 세션 메시지 전송 (tools/call 등)
|
|
||||||
transport.handleMessage(sessionId, body);
|
|
||||||
return ResponseEntity.accepted().build();
|
|
||||||
}
|
}
|
||||||
|
|
||||||
|
org.springframework.web.servlet.mvc.method.annotation.SseEmitter emitter = transport.handleCustomSse(activeSessionId, body);
|
||||||
|
|
||||||
|
return ResponseEntity.ok()
|
||||||
|
.header("Mcp-Session-Id", activeSessionId)
|
||||||
|
.body(emitter);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user