Class: MCPClient::HttpTransportBase::SseEventScanner

Inherits:
Object
  • Object
show all
Defined in:
lib/mcp_client/http_transport_base/sse_event_scanner.rb

Overview

Splits an SSE body into complete events as its bytes arrive. Line terminators are CRLF, CR or LF (SSE "Parsing an event stream"): a CR ends a line on its own, so an event it terminates is dispatched at once, and the LF of a CRLF that arrives in the next chunk is skipped. A gzip body (Streamable HTTP offers gzip on every request) is inflated as it arrives. Only a body that starts like an event stream is scanned; a JSON body never yields anything. Events are counted, terminated or not dispatched, in the same order the completed body splits into them.

Constant Summary collapse

SSE_START =

An event stream opens with a comment or a field — any field: one the client does not know is ignored, not a reason to stop reading (SSE "Parsing an event stream"). A body that opens a JSON value is never an event stream.

/\A(?::|[^\n:{\[]+:)/n
BOM =
"\xEF\xBB\xBF".b
GZIP_MAGIC =
"\x1F\x8B".b

Instance Attribute Summary collapse

Instance Method Summary collapse

Constructor Details

#initialize(max_inflated_bytes: nil) ⇒ SseEventScanner

Returns a new instance of SseEventScanner.

Parameters:

  • max_inflated_bytes (Integer, nil) (defaults to: nil) —

    bound on a gzip body's expansion, beyond which the stream is no longer scanned



30
31
32
33
34
35
36
37
38
39
40
# File 'lib/mcp_client/http_transport_base/sse_event_scanner.rb', line 30

def initialize(max_inflated_bytes: nil)
  @normalized = +''.b
  @head = +''.b
  @scanned = 0
  @after_cr = false
  @count = 0
  @sse = nil
  @inflater = nil
  @inflated = 0
  @max_inflated_bytes = max_inflated_bytes
end

Instance Attribute Details

#count ⇒ Integer (readonly)

Returns complete events seen so far.

Returns:

  • (Integer) —

    complete events seen so far



26
27
28
# File 'lib/mcp_client/http_transport_base/sse_event_scanner.rb', line 26

def count
  @count
end

Instance Method Details

#feed(chunk) {|event| ... } ⇒ void

This method returns an undefined value.

Parameters:

  • chunk (String) —

    the bytes that just arrived

Yield Parameters:

  • event (String) —

    one complete event, LF-normalized, without its terminator



45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
# File 'lib/mcp_client/http_transport_base/sse_event_scanner.rb', line 45

def feed(chunk)
  return if @sse == false

  # Bytes, not characters: the body is peer-controlled and may not be
  # text at all.
  text = decoded(chunk.b)
  return if text.nil? || @sse == false

  text = text[1..] if @after_cr && text.start_with?("\n")
  @after_cr = text.end_with?("\r")
  @normalized << text.gsub(/\r\n|\r/, "\n")
  return unless scanning?

  while (index = @normalized.index("\n\n", @scanned))
    event = @normalized[@scanned...index]
    @scanned = index + 2
    @count += 1
    # Blank lines before an event's first field dispatch nothing (SSE
    # "Parsing an event stream"), so they are not part of the event.
    yield event.sub(/\A\n+/, '').force_encoding(Encoding::UTF_8)
  end
end