From d6801b7ba0bd9e68b6d5dafb975cb8ae49919920 Mon Sep 17 00:00:00 2001 From: jade Date: Mon, 20 Jul 2026 09:36:15 +0900 Subject: [PATCH] =?UTF-8?q?fix:=20Deep=20Agent=20Builder=20=ED=98=B8?= =?UTF-8?q?=ED=99=98=EC=84=B1=EC=9D=84=20=EC=9C=84=ED=95=9C=20mcp-http=20W?= =?UTF-8?q?arm=20Pool=20=EC=84=B8=EC=85=98=20=EB=B2=84=EA=B7=B8=20?= =?UTF-8?q?=EB=B0=8F=20CORS=20=ED=97=A4=EB=8D=94=20=EC=88=98=EC=A0=95?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../dap/common/config/CorsConfig.java | 1 + .../dap/common/mcp/config/WebConfig.java | 1 + ...ustomWebMvcSseServerTransportProvider.java | 81 +++++++++++++------ .../gateway/sync/DynamicMcpController.java | 24 +++--- 4 files changed, 71 insertions(+), 36 deletions(-) diff --git a/dap-common/src/main/java/io/shinhanlife/dap/common/config/CorsConfig.java b/dap-common/src/main/java/io/shinhanlife/dap/common/config/CorsConfig.java index 6012ff3..6ebb80b 100644 --- a/dap-common/src/main/java/io/shinhanlife/dap/common/config/CorsConfig.java +++ b/dap-common/src/main/java/io/shinhanlife/dap/common/config/CorsConfig.java @@ -27,6 +27,7 @@ public class CorsConfig implements WebMvcConfigurer { .allowedOriginPatterns("*") // 외부 Agent Builder 등 모든 오리진 허용 .allowedMethods("GET", "POST", "PUT", "DELETE", "OPTIONS", "HEAD", "PATCH") // 허용할 HTTP 메서드 .allowedHeaders("*") // 모든 헤더 허용 + .exposedHeaders("Mcp-Session-Id") // MCP-HTTP 세션 아이디 노출 허용 .allowCredentials(true) // 쿠키/인증 정보 허용 .maxAge(3600); // preflight 요청 캐시 시간 (초 단위) } diff --git a/dap-common/src/main/java/io/shinhanlife/dap/common/mcp/config/WebConfig.java b/dap-common/src/main/java/io/shinhanlife/dap/common/mcp/config/WebConfig.java index a3901cb..2b9962e 100644 --- a/dap-common/src/main/java/io/shinhanlife/dap/common/mcp/config/WebConfig.java +++ b/dap-common/src/main/java/io/shinhanlife/dap/common/mcp/config/WebConfig.java @@ -48,6 +48,7 @@ public class WebConfig implements WebMvcConfigurer { .allowedOriginPatterns("*") .allowedMethods("GET", "POST", "PUT", "DELETE", "OPTIONS") .allowedHeaders("*") + .exposedHeaders("Mcp-Session-Id") .allowCredentials(true); } } \ No newline at end of file diff --git a/dap-gateway/src/main/java/io/shinhanlife/dap/biz/mcp/gateway/sync/CustomWebMvcSseServerTransportProvider.java b/dap-gateway/src/main/java/io/shinhanlife/dap/biz/mcp/gateway/sync/CustomWebMvcSseServerTransportProvider.java index ef59599..700aabc 100644 --- a/dap-gateway/src/main/java/io/shinhanlife/dap/biz/mcp/gateway/sync/CustomWebMvcSseServerTransportProvider.java +++ b/dap-gateway/src/main/java/io/shinhanlife/dap/biz/mcp/gateway/sync/CustomWebMvcSseServerTransportProvider.java @@ -49,6 +49,7 @@ public class CustomWebMvcSseServerTransportProvider implements McpServerTranspor private final String sseEndpoint; private final String messageEndpoint; private final Map sessions = new ConcurrentHashMap<>(); + private final Map customTransports = new ConcurrentHashMap<>(); private final ObjectMapper objectMapper; public CustomWebMvcSseServerTransportProvider(String sseEndpoint, String messageEndpoint, ObjectMapper objectMapper) { @@ -71,6 +72,7 @@ public class CustomWebMvcSseServerTransportProvider implements McpServerTranspor } catch (Exception ignored) {} }); sessions.clear(); + customTransports.clear(); }); } @@ -108,27 +110,38 @@ public class CustomWebMvcSseServerTransportProvider implements McpServerTranspor 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) { throw new IllegalStateException("SessionFactory not configured"); } - - SseEmitter emitter = new SseEmitter(-1L); - - CustomMcpSessionTransport sessionTransport = new CustomMcpSessionTransport(emitter, sessionId); - McpServerSession session = sessionFactory.create(sessionTransport); - sessions.put(sessionId, session); - - emitter.onCompletion(() -> sessions.remove(sessionId)); - emitter.onTimeout(() -> sessions.remove(sessionId)); - + + org.springframework.web.servlet.mvc.method.annotation.SseEmitter emitter = new org.springframework.web.servlet.mvc.method.annotation.SseEmitter(-1L); + + boolean isNew = !sessions.containsKey(sessionId); + + if (isNew) { + 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(() -> { try { - // 커스텀 클라이언트는 endpoint 이벤트를 무시할 수 있지만, 표준 호환성을 위해 전송 + // 커스텀 클라이언트는 endpoint 이벤트를 무시할 수 있지만 표준 호환성을 위해 전송 Thread.sleep(100); - emitter.send(SseEmitter.event().name("endpoint").data(messageEndpoint + "?sessionId=" + sessionId)); - - // Body로 들어온 initialize 등 즉시 처리 + emitter.send(org.springframework.web.servlet.mvc.method.annotation.SseEmitter.event().name("endpoint").data(messageEndpoint + "?sessionId=" + sessionId)); + + // Body로 들어온 메시지 즉시 비동기 처리 if (body != null && !body.trim().isEmpty()) { handleMessage(sessionId, body); } @@ -136,14 +149,14 @@ public class CustomWebMvcSseServerTransportProvider implements McpServerTranspor emitter.completeWithError(e); } }).start(); - + return emitter; } - public ResponseEntity handleMessage(String sessionId, String body) { + public org.springframework.http.ResponseEntity handleMessage(String sessionId, String body) { log.info("Received POST message for sessionId: " + sessionId + ", body: " + body); 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); @@ -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 final SseEmitter emitter; + private volatile org.springframework.web.servlet.mvc.method.annotation.SseEmitter emitter; 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.sessionId = sessionId; } + public void setEmitter(org.springframework.web.servlet.mvc.method.annotation.SseEmitter emitter) { + this.emitter = emitter; + } + @Override public Mono sendMessage(McpSchema.JSONRPCMessage message) { return Mono.fromRunnable(() -> { @@ -187,18 +208,30 @@ public class CustomWebMvcSseServerTransportProvider implements McpServerTranspor try { String json = objectMapper.writeValueAsString(message); log.info("Serialized message: " + json); - emitter.send(SseEmitter.event().name("message").data(json)); - log.info("Message successfully sent to SSE emitter"); + if (this.emitter != null) { + 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) { log.error("Error sending message to SSE emitter", e); - emitter.completeWithError(e); + if (this.emitter != null) { + this.emitter.completeWithError(e); + } } }); } @Override public Mono closeGracefully() { - return Mono.fromRunnable(emitter::complete); + return Mono.fromRunnable(() -> { + if (this.emitter != null) { + this.emitter.complete(); + } + }); } @Override diff --git a/dap-gateway/src/main/java/io/shinhanlife/dap/biz/mcp/gateway/sync/DynamicMcpController.java b/dap-gateway/src/main/java/io/shinhanlife/dap/biz/mcp/gateway/sync/DynamicMcpController.java index 1a89d8e..5cf0279 100644 --- a/dap-gateway/src/main/java/io/shinhanlife/dap/biz/mcp/gateway/sync/DynamicMcpController.java +++ b/dap-gateway/src/main/java/io/shinhanlife/dap/biz/mcp/gateway/sync/DynamicMcpController.java @@ -70,18 +70,18 @@ public class DynamicMcpController { return ResponseEntity.badRequest().body("Unknown category: " + category); } - if (sessionId == null || sessionId.isEmpty()) { - // 새 세션 생성 (initialize 요청) - String newSessionId = UUID.randomUUID().toString(); - SseEmitter emitter = transport.handleCustomSse(newSessionId, body); - - return ResponseEntity.ok() - .header("Mcp-Session-Id", newSessionId) - .body(emitter); - } else { - // 기존 세션 메시지 전송 (tools/call 등) - transport.handleMessage(sessionId, body); - return ResponseEntity.accepted().build(); + boolean isNew = (sessionId == null || sessionId.isEmpty()); + String activeSessionId = isNew ? java.util.UUID.randomUUID().toString() : sessionId; + + if (!isNew && !transport.hasSession(activeSessionId)) { + // 클라이언트가 보낸 세션 ID가 만료되었거나 존재하지 않는 경우 (Warm Pool 스펙: 404 Not Found 반환) + return ResponseEntity.notFound().build(); } + + org.springframework.web.servlet.mvc.method.annotation.SseEmitter emitter = transport.handleCustomSse(activeSessionId, body); + + return ResponseEntity.ok() + .header("Mcp-Session-Id", activeSessionId) + .body(emitter); } }