Class: ModelContextProtocol::Server::StreamableHttpTransport::SessionMessageQueue

Inherits:
Object
  • Object
show all
Defined in:
lib/model_context_protocol/server/streamable_http_transport/session_message_queue.rb

Constant Summary collapse

QUEUE_KEY_PREFIX =
"session_messages:"
DEFAULT_TTL =

1 hour

3600
MAX_MESSAGES =
1000

Instance Method Summary collapse

Constructor Details

#initialize(redis_client, session_id, ttl: DEFAULT_TTL) ⇒ SessionMessageQueue

Returns a new instance of SessionMessageQueue.



10
11
12
13
14
15
# File 'lib/model_context_protocol/server/streamable_http_transport/session_message_queue.rb', line 10

def initialize(redis_client, session_id, ttl: DEFAULT_TTL)
  @redis = redis_client
  @session_id = session_id
  @queue_key = "#{QUEUE_KEY_PREFIX}#{session_id}"
  @ttl = ttl
end

Instance Method Details

#has_messages?Boolean

Returns:



43
44
45
46
47
# File 'lib/model_context_protocol/server/streamable_http_transport/session_message_queue.rb', line 43

def has_messages?
  @redis.exists(@queue_key) > 0
rescue
  false
end

#poll_messagesObject



27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
# File 'lib/model_context_protocol/server/streamable_http_transport/session_message_queue.rb', line 27

def poll_messages
  lua_script = "    local messages = redis.call('lrange', KEYS[1], 0, -1)\n    if #messages > 0 then\n      redis.call('del', KEYS[1])\n    end\n    return messages\n  LUA\n\n  messages = @redis.eval(lua_script, keys: [@queue_key])\n  return [] unless messages && !messages.empty?\n  messages.reverse.map { |json| deserialize_message(json) }\nrescue\n  []\nend\n"

#push_message(message) ⇒ Object



17
18
19
20
21
22
23
24
25
# File 'lib/model_context_protocol/server/streamable_http_transport/session_message_queue.rb', line 17

def push_message(message)
  message_json = serialize_message(message)

  @redis.multi do |multi|
    multi.lpush(@queue_key, message_json)
    multi.expire(@queue_key, @ttl)
    multi.ltrim(@queue_key, 0, MAX_MESSAGES - 1)
  end
end