Class: ModelContextProtocol::Server::StreamableHttpTransport::ServerRequestStore

Inherits:
Object
  • Object
show all
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

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

Parameters:

  • session_id (String)

    the session identifier

Returns:

  • (Array<String>)

    list of cleaned up request IDs



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

Parameters:

  • timeout_seconds (Integer)

    timeout in seconds

Returns:

  • (Array<Hash>)

    array of expired request info with request_id and session_id



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

Parameters:

  • request_id (String)

    the unique JSON-RPC request identifier

Returns:

  • (Hash, nil)

    request information or nil if not found



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)

Parameters:

  • request_id (String)

    the unique JSON-RPC request identifier

Returns:

  • (Boolean)

    true if request was pending, false if not found



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

Parameters:

  • request_id (String)

    the unique JSON-RPC request identifier

Returns:

  • (Boolean)

    true if the request is pending, false otherwise



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

Parameters:

  • request_id (String)

    the unique JSON-RPC request identifier

  • session_id (String) (defaults to: nil)

    the session identifier (can be nil for sessionless requests)

  • type (Symbol) (defaults to: :ping)

    the type of request (e.g., :ping)



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)

Parameters:

  • request_id (String)

    the unique JSON-RPC request identifier



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