Class: ModelContextProtocol::Server::StreamableHttpTransport::ServerRequestStore
- Inherits:
-
Object
- Object
- ModelContextProtocol::Server::StreamableHttpTransport::ServerRequestStore
- Defined in:
- lib/model_context_protocol/server/streamable_http_transport/server_request_store.rb
Overview
Redis-based distributed storage for tracking server-initiated requests and their response status. This store is used by StreamableHttpTransport to manage outgoing request lifecycle (like pings) across multiple server instances and handle timeouts in a distributed environment.
Constant Summary collapse
- REQUEST_KEY_PREFIX =
"server_request:pending:"- SESSION_KEY_PREFIX =
"server_request:session:"- DEFAULT_TTL =
1 minute TTL for request entries
60
Instance Method Summary collapse
-
#cleanup_session_requests(session_id) ⇒ Array<String>
Clean up all server requests associated with a session This is typically called when a session is terminated.
-
#get_expired_requests(timeout_seconds) ⇒ Array<Hash>
Find requests that have exceeded the specified timeout.
-
#get_request(request_id) ⇒ Hash?
Get information about a specific pending request.
-
#initialize(redis_client, server_instance, ttl: DEFAULT_TTL) ⇒ ServerRequestStore
constructor
A new instance of ServerRequestStore.
-
#mark_completed(request_id) ⇒ Boolean
Mark a server-initiated request as completed (response received).
-
#pending?(request_id) ⇒ Boolean
Check if a server-initiated request is still pending.
-
#register_request(request_id, session_id = nil, type: :ping) ⇒ void
Register a new server-initiated request with its associated session.
-
#unregister_request(request_id) ⇒ void
Unregister a request (typically called when request completes or times out).
Constructor Details
#initialize(redis_client, server_instance, ttl: DEFAULT_TTL) ⇒ ServerRequestStore
Returns a new instance of ServerRequestStore.
13 14 15 16 17 |
# File 'lib/model_context_protocol/server/streamable_http_transport/server_request_store.rb', line 13 def initialize(redis_client, server_instance, ttl: DEFAULT_TTL) @redis = redis_client @server_instance = server_instance @ttl = ttl end |
Instance Method Details
#cleanup_session_requests(session_id) ⇒ Array<String>
Clean up all server requests associated with a session This is typically called when a session is terminated
145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 |
# File 'lib/model_context_protocol/server/streamable_http_transport/server_request_store.rb', line 145 def cleanup_session_requests(session_id) pattern = "#{SESSION_KEY_PREFIX}#{session_id}:*" request_keys = @redis.keys(pattern) return [] if request_keys.empty? # Extract request IDs from the keys request_ids = request_keys.map do |key| key.sub("#{SESSION_KEY_PREFIX}#{session_id}:", "") end # Delete all related keys all_keys = [] request_ids.each do |request_id| all_keys << "#{REQUEST_KEY_PREFIX}#{request_id}" end all_keys.concat(request_keys) @redis.del(*all_keys) unless all_keys.empty? request_ids end |
#get_expired_requests(timeout_seconds) ⇒ Array<Hash>
Find requests that have exceeded the specified timeout
79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 |
# File 'lib/model_context_protocol/server/streamable_http_transport/server_request_store.rb', line 79 def get_expired_requests(timeout_seconds) current_time = Time.now.to_f expired_requests = [] # Get all pending request keys request_keys = @redis.keys("#{REQUEST_KEY_PREFIX}*") return expired_requests if request_keys.empty? # Get all request data in batch request_values = @redis.mget(request_keys) request_keys.each_with_index do |key, index| next unless request_values[index] begin request_data = JSON.parse(request_values[index]) created_at = request_data["created_at"] if created_at && (current_time - created_at) > timeout_seconds request_id = key.sub(REQUEST_KEY_PREFIX, "") expired_requests << { request_id: request_id, session_id: request_data["session_id"], type: request_data["type"], age: current_time - created_at } end rescue JSON::ParserError # Skip malformed entries next end end expired_requests end |
#get_request(request_id) ⇒ Hash?
Get information about a specific pending request
68 69 70 71 72 73 |
# File 'lib/model_context_protocol/server/streamable_http_transport/server_request_store.rb', line 68 def get_request(request_id) data = @redis.get("#{REQUEST_KEY_PREFIX}#{request_id}") data ? JSON.parse(data) : nil rescue JSON::ParserError nil end |
#mark_completed(request_id) ⇒ Boolean
Mark a server-initiated request as completed (response received)
48 49 50 51 52 53 54 |
# File 'lib/model_context_protocol/server/streamable_http_transport/server_request_store.rb', line 48 def mark_completed(request_id) request_data = @redis.get("#{REQUEST_KEY_PREFIX}#{request_id}") return false unless request_data unregister_request(request_id) true end |
#pending?(request_id) ⇒ Boolean
Check if a server-initiated request is still pending
60 61 62 |
# File 'lib/model_context_protocol/server/streamable_http_transport/server_request_store.rb', line 60 def pending?(request_id) @redis.exists("#{REQUEST_KEY_PREFIX}#{request_id}") == 1 end |
#register_request(request_id, session_id = nil, type: :ping) ⇒ void
This method returns an undefined value.
Register a new server-initiated request with its associated session
25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 |
# File 'lib/model_context_protocol/server/streamable_http_transport/server_request_store.rb', line 25 def register_request(request_id, session_id = nil, type: :ping) request_data = { session_id: session_id, server_instance: @server_instance, type: type.to_s, created_at: Time.now.to_f } @redis.multi do |multi| multi.set("#{REQUEST_KEY_PREFIX}#{request_id}", request_data.to_json, ex: @ttl) if session_id multi.set("#{SESSION_KEY_PREFIX}#{session_id}:#{request_id}", true, ex: @ttl) end end end |
#unregister_request(request_id) ⇒ void
This method returns an undefined value.
Unregister a request (typically called when request completes or times out)
119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 |
# File 'lib/model_context_protocol/server/streamable_http_transport/server_request_store.rb', line 119 def unregister_request(request_id) request_data = @redis.get("#{REQUEST_KEY_PREFIX}#{request_id}") keys_to_delete = ["#{REQUEST_KEY_PREFIX}#{request_id}"] if request_data begin data = JSON.parse(request_data) session_id = data["session_id"] if session_id keys_to_delete << "#{SESSION_KEY_PREFIX}#{session_id}:#{request_id}" end rescue JSON::ParserError nil end end @redis.del(*keys_to_delete) unless keys_to_delete.empty? end |