Class: ModelContextProtocol::Server::StreamableHttpTransport::SessionMessageQueue
- Inherits:
-
Object
- Object
- ModelContextProtocol::Server::StreamableHttpTransport::SessionMessageQueue
- 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
- #has_messages? ⇒ Boolean
-
#initialize(redis_client, session_id, ttl: DEFAULT_TTL) ⇒ SessionMessageQueue
constructor
A new instance of SessionMessageQueue.
- #poll_messages ⇒ Object
- #push_message(message) ⇒ Object
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
43 44 45 46 47 |
# File 'lib/model_context_protocol/server/streamable_http_transport/session_message_queue.rb', line 43 def @redis.exists(@queue_key) > 0 rescue false end |
#poll_messages ⇒ Object
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 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 () = () @redis.multi do |multi| multi.lpush(@queue_key, ) multi.expire(@queue_key, @ttl) multi.ltrim(@queue_key, 0, MAX_MESSAGES - 1) end end |