Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -4,26 +4,27 @@

package io.modelcontextprotocol.spec;

import java.time.Duration;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicLong;
import java.util.concurrent.atomic.AtomicReference;

import io.modelcontextprotocol.common.McpTransportContext;
import io.modelcontextprotocol.json.TypeRef;
import io.modelcontextprotocol.server.McpAsyncServerExchange;
import io.modelcontextprotocol.server.McpInitRequestHandler;
import io.modelcontextprotocol.server.McpNotificationHandler;
import io.modelcontextprotocol.server.McpRequestHandler;
import io.modelcontextprotocol.json.TypeRef;
import io.modelcontextprotocol.util.Assert;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import reactor.core.publisher.Mono;
import reactor.core.publisher.MonoSink;
import reactor.core.publisher.Sinks;

import java.time.Duration;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicLong;
import java.util.concurrent.atomic.AtomicReference;

/**
* Represents a Model Context Protocol (MCP) session on the server side. It manages
* bidirectional JSON-RPC communication with the client.
Expand All @@ -36,7 +37,9 @@ public class McpServerSession implements McpLoggableSession {

private final String id;

/** Duration to wait for request responses before timing out */
/**
* Duration to wait for request responses before timing out
*/
private final Duration requestTimeout;

private final AtomicLong requestCounter = new AtomicLong(0);
Expand Down Expand Up @@ -65,6 +68,8 @@ public class McpServerSession implements McpLoggableSession {

private volatile McpSchema.LoggingLevel minLoggingLevel = McpSchema.LoggingLevel.INFO;

private volatile AtomicBoolean closed = new AtomicBoolean(false);

/**
* Creates a new server session with the given parameters and the transport to use.
* @param id session id
Expand Down Expand Up @@ -345,14 +350,23 @@ private MethodNotFoundError getMethodNotFoundError(String method) {

@Override
public Mono<Void> closeGracefully() {
// TODO: clear pendingResponses and emit errors?
return this.transport.closeGracefully();
if (this.closed.compareAndSet(false, true)) {
this.pendingResponses.forEach((id, response) -> response.error(new RuntimeException("Session closed")));
this.pendingResponses.clear();
return this.transport.closeGracefully();
}
else {
return Mono.empty();
}
}

@Override
public void close() {
// TODO: clear pendingResponses and emit errors?
this.transport.close();
if (this.closed.compareAndSet(false, true)) {
this.pendingResponses.forEach((id, response) -> response.error(new RuntimeException("Session closed")));
this.pendingResponses.clear();
this.transport.close();
}
}

/**
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -263,10 +263,13 @@ private ServerResponse handleSseConnection(ServerRequest request) {
logger.debug("Creating new SSE connection for session: {}", sessionId);
sseBuilder.onComplete(() -> {
logger.debug("SSE connection completed for session: {}", sessionId);
// explicitly close the session when the SSE connection is completed
session.close();
sessions.remove(sessionId);
});
sseBuilder.onTimeout(() -> {
logger.debug("SSE connection timed out for session: {}", sessionId);
session.close();
sessions.remove(sessionId);
});
this.sessions.put(sessionId, session);
Expand Down Expand Up @@ -383,6 +386,12 @@ public Mono<Void> sendMessage(McpSchema.JSONRPCMessage message) {
String jsonText = jsonMapper.writeValueAsString(message);
sseBuilder.event(MESSAGE_EVENT_TYPE).data(jsonText);
}
catch (IOException e) {
if (logger.isDebugEnabled()) {
logger.debug("Failed to send message: {}", e.getMessage());
}
sseBuilder.error(e);
}
catch (Exception e) {
logger.error("Failed to send message: {}", e.getMessage());
sseBuilder.error(e);
Expand Down