Class: Raptor::Http1

Inherits:
Object
  • Object
show all
Defined in:
lib/raptor/http1.rb,
sig/generated/raptor/http1.rbs

Overview

Parses HTTP/1.x requests and dispatches them to the Rack application. Coordinates with the Ractor pool for parsing and with the reactor for requests that need more data before they can be handled.

Constant Summary collapse

BODY_BUFFER_THRESHOLD =

Returns:

  • (Object)
256 * 1024
CHUNKED_WRITE_THRESHOLD =

Returns:

  • (Object)
512 * 1024
FILE_CHUNK_SIZE =

Returns:

  • (Object)
64 * 1024
READ_BUFFER_SIZE =

Returns:

  • (Object)
64 * 1024
RESPONSE_BUFFER_CAPACITY =

Returns:

  • (Object)
4 * 1024
KEEPALIVE_READ_TIMEOUT =

Returns:

  • (::Float)
0.001
MAX_KEEPALIVE_REQUESTS =

Returns:

  • (::Integer)
100
HTTP_10 =

Returns:

  • (::String)
"HTTP/1.0"
HTTP_11 =

Returns:

  • (::String)
"HTTP/1.1"
STATUS_LINE_CACHE_10 =

Returns:

  • (Object)
Hash.new do |h, status|
  reason = Rack::Utils::HTTP_STATUS_CODES[status]
  h[status] = "HTTP/1.0 #{status}#{reason ? " #{reason}" : ""}\r\n".freeze
end
STATUS_LINE_CACHE_11 =

Returns:

  • (Object)
Hash.new do |h, status|
  reason = Rack::Utils::HTTP_STATUS_CODES[status]
  h[status] = "HTTP/1.1 #{status}#{reason ? " #{reason}" : ""}\r\n".freeze
end
CONTINUE_RESPONSE =

Returns:

  • (::String)
"HTTP/1.1 100 Continue\r\n\r\n"
BAD_REQUEST_RESPONSE =

Returns:

  • (::String)
"HTTP/1.1 400 Bad Request\r\nContent-Length: 0\r\nConnection: close\r\n\r\n"
CONTENT_TOO_LARGE_RESPONSE =

Returns:

  • (::String)
"HTTP/1.1 413 Content Too Large\r\nContent-Length: 0\r\nConnection: close\r\n\r\n"
INTERNAL_SERVER_ERROR_RESPONSE =

Returns:

  • (::String)
"HTTP/1.1 500 Internal Server Error\r\nContent-Length: 0\r\nConnection: close\r\n\r\n"
CONNECTION_CLOSE =

Returns:

  • (::String)
"close"
CONNECTION_KEEPALIVE =

Returns:

  • (::String)
"keep-alive"
EXPECT_100_CONTINUE =

Returns:

  • (::String)
"100-continue"
TRANSFER_ENCODING_CHUNKED =

Returns:

  • (::String)
"chunked"
CHUNKED_TERMINATOR =

Returns:

  • (::String)
"0\r\n\r\n"
HTTP_CONNECTION =

Returns:

  • (::String)
"HTTP_CONNECTION"
HTTP_EXPECT =

Returns:

  • (::String)
"HTTP_EXPECT"
HTTP_TRANSFER_ENCODING =

Returns:

  • (::String)
"HTTP_TRANSFER_ENCODING"
RACK_HEADER_PREFIX =

Returns:

  • (::String)
"rack."
RACK_HIJACKED =

Returns:

  • (::String)
"rack.hijacked"
RACK_HIJACK_IO =

Returns:

  • (::String)
"rack.hijack_io"

Class Method Summary collapse

Instance Method Summary collapse

Constructor Details

#initialize(app, server_port, connection_options: {}, http1_options: {}, access_log_io: nil, on_error: nil) ⇒ Http1

Creates a new Http1 handler.

RBS:

  • (^(Hash[String, untyped]) -> [Integer, Hash[String, String | Array[String]], untyped] app, Integer server_port, ?connection_options: Hash[Symbol, untyped], ?http1_options: Hash[Symbol, untyped], ?access_log_io: IO?, ?on_error: ^(Hash[String, untyped]?, Exception) -> void | nil) -> void

Parameters:

  • app (#call)

    the Rack application to dispatch complete requests to

  • server_port (Integer)

    port number used to populate SERVER_PORT in the Rack env

  • connection_options (Hash) (defaults to: {})

    per-connection settings shared across protocols

  • http1_options (Hash) (defaults to: {})

    HTTP/1.1-specific settings

  • access_log_io (IO, nil) (defaults to: nil)

    IO to write Common Log Format access entries to, or nil to disable

  • on_error (#call, nil) (defaults to: nil)

    callback invoked with (env, exception) when the Rack app raises

  • connection_options: (Hash[Symbol, untyped]) (defaults to: {})
  • http1_options: (Hash[Symbol, untyped]) (defaults to: {})
  • access_log_io: (IO, nil) (defaults to: nil)
  • on_error: (^(Hash[String, untyped]?, Exception) -> void, nil) (defaults to: nil)

Options Hash (connection_options:):

  • :write_timeout (Integer)

    per-write socket timeout in seconds

  • :max_body_size (Integer, nil)

    maximum request body size in bytes

  • :body_spool_threshold (Integer, nil)

    spool bodies larger than this to a tempfile

Options Hash (http1_options:):

  • :max_keepalive_requests (Integer)

    maximum requests per HTTP/1.1 keep-alive connection



190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
# File 'lib/raptor/http1.rb', line 190

def initialize(app, server_port, connection_options: {}, http1_options: {}, access_log_io: nil, on_error: nil)
  @app = app
  @server_port = server_port
  @server_port_string = server_port.to_s.freeze
  @write_timeout = connection_options[:write_timeout] || Http::WRITE_TIMEOUT
  @max_body_size = connection_options[:max_body_size]
  @body_spool_threshold = connection_options[:body_spool_threshold]
  @max_keepalive_requests = http1_options[:max_keepalive_requests] || MAX_KEEPALIVE_REQUESTS
  @access_log_io = access_log_io
  @on_error = on_error
  @running = AtomicBoolean.new(true)
  @env_template = {
    Rack::RACK_VERSION => Rack::VERSION,
    Rack::RACK_IS_HIJACK => true,
    Rack::SCRIPT_NAME => "",
    Rack::QUERY_STRING => "",
    Http::SERVER_SOFTWARE => Http::SERVER_SOFTWARE_VALUE
  }.freeze
end

Class Method Details

.invalid_host?(env) ⇒ Boolean

Returns true when an HTTP/1.1 request lacks a valid Host header per RFC 9112 section 3.2, where a valid value is a non-empty single-value line.

RBS:

  • (Hash[String, untyped] env) -> bool

Parameters:

  • env (Hash)

    the Rack environment after header parsing

Returns:

  • (Boolean)


66
67
68
69
70
71
# File 'lib/raptor/http1.rb', line 66

def self.invalid_host?(env)
  return false unless env[Rack::SERVER_PROTOCOL] == HTTP_11

  http_host = env[Rack::HTTP_HOST]
  !http_host || http_host.empty? || http_host.include?(",")
end

.parse(data, env_template, max_body_size) ⇒ Hash

Advances an HTTP/1.x request parse from the state hash's buffered bytes. A complete well-formed request returns with :complete set plus populated :env and :body; malformed or oversized input flips :malformed or :too_large; incomplete input returns the state so the reactor can wait for more.

RBS:

  • (Hash[Symbol, untyped] data, Hash[String, untyped] env_template, Integer? max_body_size) -> Hash[Symbol, untyped]

Parameters:

  • data (Hash)

    the current parse state

  • env_template (Hash)

    the Rack env template to seed the request with

  • max_body_size (Integer, nil)

    byte limit for the request body, or nil for no limit

Returns:

  • (Hash)

    the updated parse state, made shareable for cross-Ractor return



112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
# File 'lib/raptor/http1.rb', line 112

def self.parse(data, env_template, max_body_size)
  parser = Raptor::HttpParser.new
  env = env_template.dup
  nread = begin
    parser.execute(env, data[:buffer], 0)
  rescue Raptor::HttpParserError
    return Ractor.make_shareable(data.merge(complete: true, malformed: true))
  end
  parse_data = if data[:parse_data]
    data[:parse_data].dup
  else
    { parse_count: 0, content_length: parser.content_length }
  end
  parse_data[:parse_count] += 1

  message = if parser.finished?
    if invalid_host?(env) || request_smuggling?(env)
      data.merge(env: env, body: nil, parse_data: parse_data, complete: true, malformed: true)
    elsif parser.has_body?
      body_buffer = data[:buffer].byteslice(nread..-1) || ""

      if max_body_size && parser.content_length > max_body_size
        data.merge(env: env, body: nil, parse_data: parse_data, complete: true, too_large: true)
      elsif parser.chunked?
        decoded_body, chunked_state = HttpParser.decode_chunked(body_buffer, max_body_size)

        case chunked_state
        when :complete
          env.delete(HTTP_TRANSFER_ENCODING)
          data.merge(env: env, body: decoded_body, parse_data: parse_data, complete: true)
        when :too_large
          data.merge(env: env, body: nil, parse_data: parse_data, complete: true, too_large: true)
        when :malformed
          data.merge(env: env, body: nil, parse_data: parse_data, complete: true, malformed: true)
        else
          data.merge(env: env, parse_data: parse_data)
        end
      elsif parser.content_length > body_buffer.bytesize
        data.merge(env: env, parse_data: parse_data)
      else
        data.merge(env: env, body: body_buffer, parse_data: parse_data, complete: true)
      end
    else
      data.merge(env: env, body: nil, parse_data: parse_data, complete: true)
    end
  else
    data.merge(env: env, parse_data: parse_data)
  end
  Ractor.make_shareable(message)
end

.request_smuggling?(env) ⇒ Boolean

Returns true when the message framing shows a request-smuggling vector per RFC 9112 section 6.3: a Transfer-Encoding where chunked is missing, not the final encoding, or duplicated; a Transfer-Encoding paired with a Content-Length; or a Content-Length containing any non-digit character.

RBS:

  • (Hash[String, untyped] env) -> bool

Parameters:

  • env (Hash)

    the Rack environment after header parsing

Returns:

  • (Boolean)


83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
# File 'lib/raptor/http1.rb', line 83

def self.request_smuggling?(env)
  transfer_encoding = env[HTTP_TRANSFER_ENCODING]
  content_length = env[Http::CONTENT_LENGTH]

  if transfer_encoding
    return true if content_length

    encodings = transfer_encoding.downcase.split(",").map(&:strip)
    return true if encodings.last != TRANSFER_ENCODING_CHUNKED
    return true if encodings.count(TRANSFER_ENCODING_CHUNKED) > 1
  elsif content_length
    return true if content_length.match?(/[^\d]/)
  end

  false
end

Instance Method Details

#build_rack_env(env, parse_data, body, socket, remote_addr: Server::DEFAULT_REMOTE_ADDR, url_scheme: Server::HTTP_SCHEME) ⇒ Hash

Builds a Rack environment hash from parsed HTTP request data.

RBS:

  • (Hash[String, untyped] env, Hash[Symbol, untyped] parse_data, String? body, TCPSocket socket, ?remote_addr: String, ?url_scheme: String) -> Hash[String, untyped]

Parameters:

  • env (Hash)

    partial env hash from the HTTP parser

  • parse_data (Hash)

    metadata from the parsing pass, including content_length

  • body (String, nil)

    decoded request body, or nil if no body

  • socket (TCPSocket)

    the client socket, used for hijack support

  • remote_addr (String) (defaults to: Server::DEFAULT_REMOTE_ADDR)

    client IP address

  • url_scheme (String) (defaults to: Server::HTTP_SCHEME)

    "http" or "https"

  • remote_addr: (String) (defaults to: Server::DEFAULT_REMOTE_ADDR)
  • url_scheme: (String) (defaults to: Server::HTTP_SCHEME)

Returns:

  • (Hash)

    fully populated Rack environment hash



728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
# File 'lib/raptor/http1.rb', line 728

def build_rack_env(env, parse_data, body, socket, remote_addr: Server::DEFAULT_REMOTE_ADDR, url_scheme: Server::HTTP_SCHEME)
  env[Rack::RACK_INPUT] = build_rack_input(body)
  env[Rack::RACK_ERRORS] = $stderr
  env[Rack::RACK_RESPONSE_FINISHED] = []
  env[Rack::RACK_HIJACK] = proc do
    env[RACK_HIJACKED] = true
    env[RACK_HIJACK_IO] = socket
    socket
  end
  env[Rack::RACK_EARLY_HINTS] = proc do |hints|
    send_early_hints(socket, hints) rescue nil
  end

  unless env.key?(Rack::PATH_INFO)
    request_uri = env[Http::REQUEST_URI]
    scheme_end = request_uri&.index("://")
    if scheme_end
      authority_end = request_uri.index("/", scheme_end + 3) || request_uri.bytesize
      path_and_query = request_uri.byteslice(authority_end..-1) || ""
      if query_delim = path_and_query.index("?")
        env[Rack::PATH_INFO] = query_delim.zero? ? "/" : path_and_query.byteslice(0, query_delim)
        env[Rack::QUERY_STRING] = path_and_query.byteslice(query_delim + 1..-1)
      else
        env[Rack::PATH_INFO] = path_and_query.empty? ? "/" : path_and_query
      end
    else
      env[Rack::PATH_INFO] = ""
    end
  end

  if (content_length = parse_data[:content_length]).positive?
    env[Http::CONTENT_LENGTH] = content_length.to_s
  end

  env[Http::REMOTE_ADDR] = remote_addr
  env[Http::HTTP_VERSION] = env[Rack::SERVER_PROTOCOL]

  behind_tls_proxy = (url_scheme == Server::HTTP_SCHEME) && forwarded_https?(env)
  env[Rack::RACK_URL_SCHEME] = behind_tls_proxy ? Server::HTTPS_SCHEME : url_scheme
  default_port = behind_tls_proxy ? "443" : @server_port_string

  http_host = env[Rack::HTTP_HOST]
  host = nil
  port = nil
  if http_host && !http_host.empty?
    if http_host.start_with?("[")
      bracket_end = http_host.index("]")
      if bracket_end
        host = http_host.byteslice(1, bracket_end - 1)
        port_colon = http_host.index(":", bracket_end + 1)
        port = port_colon && http_host.byteslice(port_colon + 1, http_host.bytesize - port_colon - 1)
      end
    else
      colon = http_host.index(":")
      if colon
        host = http_host.byteslice(0, colon)
        port = http_host.byteslice(colon + 1, http_host.bytesize - colon - 1)
      else
        host = http_host
      end
    end
  end
  env[Rack::SERVER_NAME] ||= host || Server::DEFAULT_SERVER_NAME
  env[Rack::SERVER_PORT] ||= port || default_port

  env
end

#build_rack_input(body) ⇒ IO

Builds the rack.input IO object for the request body. Returns an in-memory StringIO for bodies up to the spool threshold, or a Tempfile for larger bodies to bound per-worker memory.

RBS:

  • (String? body) -> IO

Parameters:

  • body (String, nil)

    decoded request body

Returns:

  • (IO)

    an IO-like object positioned at the start of the body



804
805
806
807
808
809
810
811
812
813
814
# File 'lib/raptor/http1.rb', line 804

def build_rack_input(body)
  if body && @body_spool_threshold && body.bytesize > @body_spool_threshold
    tempfile = Tempfile.new("raptor-body")
    tempfile.binmode
    tempfile.write(body)
    tempfile.rewind
    tempfile
  else
    (body ? StringIO.new(body) : StringIO.new).set_encoding(Encoding::ASCII_8BIT)
  end
end

#build_status_line(http_version, status) ⇒ String

Returns the HTTP status line for status.

RBS:

  • (String http_version, Integer status) -> String

Parameters:

  • http_version (String)

    "HTTP/1.1" or "HTTP/1.0"

  • status (Integer)

    HTTP status code

Returns:

  • (String)

    the status line including trailing CRLF



984
985
986
987
988
989
990
# File 'lib/raptor/http1.rb', line 984

def build_status_line(http_version, status)
  cache = http_version == HTTP_11 ? STATUS_LINE_CACHE_11 : STATUS_LINE_CACHE_10
  response = (Thread.current[:raptor_response_buffer] ||= String.new(capacity: RESPONSE_BUFFER_CAPACITY))
  response.clear
  response << cache[status]
  response
end

#calculate_content_length(body) ⇒ Integer?

Returns the byte length of the body when it can be determined upfront (array or file), otherwise nil.

RBS:

  • (untyped body) -> Integer?

Parameters:

  • body (Object)

    the response body

Returns:

  • (Integer, nil)

    the byte length, or nil if it cannot be determined



1092
1093
1094
1095
1096
1097
1098
1099
1100
1101
1102
1103
# File 'lib/raptor/http1.rb', line 1092

def calculate_content_length(body)
  if body.respond_to?(:to_ary)
    array = body.to_ary
    return unless array.is_a?(Array)

    array.sum { |chunk| chunk.is_a?(String) ? chunk.bytesize : 0 }
  elsif body.respond_to?(:to_path) && (path = body.to_path) && File.readable?(path)
    File.size(path)
  else
    nil
  end
end

#call_response_finished(env, status, headers, error) ⇒ void

This method returns an undefined value.

Calls every rack.response_finished callback in reverse registration order, rescuing any that raise.

RBS:

  • (Hash[String, untyped] env, Integer? status, Hash[String, String | Array[String]]? headers, Exception? error) -> void

Parameters:

  • env (Hash, nil)

    the Rack environment

  • status (Integer, nil)

    the response status code

  • headers (Hash, nil)

    the response headers

  • error (Exception, nil)

    any error raised during processing, or nil on success



1251
1252
1253
1254
1255
1256
1257
# File 'lib/raptor/http1.rb', line 1251

def call_response_finished(env, status, headers, error)
  return unless env && env[Rack::RACK_RESPONSE_FINISHED].is_a?(Array)

  env[Rack::RACK_RESPONSE_FINISHED].reverse_each do |callable|
    callable.call(env, status, headers, error) rescue nil
  end
end

#cork_socket(socket) ⇒ void

This method returns an undefined value.

Enables TCP_CORK on the socket to batch outgoing packets into fewer segments. Linux-only; a no-op elsewhere.

RBS:

  • (TCPSocket socket) -> void

Parameters:

  • socket (TCPSocket)

    the socket to cork



1294
1295
1296
# File 'lib/raptor/http1.rb', line 1294

def cork_socket(socket)
  socket.setsockopt(Socket::IPPROTO_TCP, Socket::TCP_CORK, 1) if socket.is_a?(TCPSocket)
end

#eager_accept(socket, id, reactor, thread_pool, remote_addr, url_scheme) ⇒ void

This method returns an undefined value.

Eagerly reads and parses the first request on a freshly accepted connection on the server thread, dispatching directly to the thread pool when complete. Falls back to the reactor when more data is needed.

RBS:

  • (TCPSocket socket, Integer id, Reactor reactor, AtomicThreadPool thread_pool, String remote_addr, String url_scheme) -> void

Parameters:

  • socket (TCPSocket)

    the freshly accepted client socket

  • id (Integer)

    unique client identifier

  • reactor (Reactor)

    the reactor for fallback registration

  • thread_pool (AtomicThreadPool)

    thread pool for application processing

  • remote_addr (String)

    client IP address

  • url_scheme (String)

    "http" or "https"



269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
# File 'lib/raptor/http1.rb', line 269

def eager_accept(socket, id, reactor, thread_pool, remote_addr, url_scheme)
  begin
    buffer = read_into_thread_buffer(socket)
  rescue IO::WaitReadable
    reactor.add(id: id, socket: socket, remote_addr: remote_addr, url_scheme: url_scheme)
    return
  rescue EOFError, IOError
    socket.close rescue nil
    return
  end

  env, parse_data, nread, parser = begin
    parse_next_request(buffer)
  rescue HttpParserError
    reject_malformed(socket)
    return
  end

  if !parser.finished?
    fallback_to_reactor(socket, id, buffer, env, parse_data, reactor, 0, remote_addr, url_scheme, persisted: false)
    return
  elsif Http1.invalid_host?(env) || Http1.request_smuggling?(env)
    reject_malformed(socket)
    return
  elsif parser.has_body? && @max_body_size && parser.content_length > @max_body_size
    reject_oversized(socket)
    return
  end

  body = extract_body(buffer, env, parser, nread, decode_chunked: true)
  case body
  when :incomplete
    fallback_to_reactor(socket, id, buffer, env, parse_data, reactor, 0, remote_addr, url_scheme, persisted: false)
    return
  when :too_large
    reject_oversized(socket)
    return
  when :malformed
    reject_malformed(socket)
    return
  end

  thread_pool << proc do
    process_client(socket, id, env, parse_data, body, reactor, thread_pool, 1, remote_addr, url_scheme)
  end
end

#eager_keepalive(socket, id, reactor, thread_pool, request_count, remote_addr, url_scheme) ⇒ void

This method returns an undefined value.

Reads and processes subsequent requests inline on a kept-alive connection. Falls back to the reactor when no data arrives within the timeout, the thread pool is saturated, or the request is incomplete.

RBS:

  • (TCPSocket socket, Integer id, Reactor reactor, AtomicThreadPool thread_pool, Integer request_count, String remote_addr, String url_scheme) -> void

Parameters:

  • socket (TCPSocket)

    the client socket

  • id (Integer)

    unique client identifier

  • reactor (Reactor)

    the reactor for fallback registration

  • thread_pool (AtomicThreadPool)

    thread pool for deprioritization

  • request_count (Integer)

    number of requests handled on this connection

  • remote_addr (String)

    client IP address

  • url_scheme (String)

    "http" or "https"



587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
# File 'lib/raptor/http1.rb', line 587

def eager_keepalive(socket, id, reactor, thread_pool, request_count, remote_addr, url_scheme)
  loop do
    unless @running.true?
      socket.close rescue nil
      return
    end

    unless socket.wait_readable(KEEPALIVE_READ_TIMEOUT)
      reactor.persist(socket, id, request_count, remote_addr: remote_addr, url_scheme: url_scheme)
      return
    end

    begin
      buffer = read_into_thread_buffer(socket)
    rescue IO::WaitReadable
      reactor.persist(socket, id, request_count, remote_addr: remote_addr, url_scheme: url_scheme)
      return
    rescue EOFError
      socket.close rescue nil
      return
    end

    env, parse_data, nread, parser = begin
      parse_next_request(buffer)
    rescue HttpParserError
      reject_malformed(socket)
      return
    end

    if !parser.finished?
      fallback_to_reactor(socket, id, buffer, env, parse_data, reactor, request_count, remote_addr, url_scheme)
      return
    end

    body = extract_body(buffer, env, parser, nread, decode_chunked: false)
    if body == :incomplete
      fallback_to_reactor(socket, id, buffer, env, parse_data, reactor, request_count, remote_addr, url_scheme)
      return
    end

    request_count += 1

    if thread_pool.queue_size >= thread_pool.size
      thread_pool << proc do
        process_client(
          socket,
          id,
          env,
          parse_data,
          body,
          reactor,
          thread_pool,
          request_count,
          remote_addr,
          url_scheme
        )
      end
      return
    end

    keep_alive = process_request(
      socket,
      env,
      parse_data,
      body,
      request_count,
      remote_addr,
      url_scheme
    )
    return unless keep_alive
  end
end

#expects_100_continue?(env) ⇒ Boolean

Returns true if the request expects a 100 Continue response per RFC 7231 section 5.1.1.

RBS:

  • (Hash[String, untyped] env) -> bool

Parameters:

  • env (Hash)

    the parsed Rack environment (possibly incomplete)

Returns:

  • (Boolean)


446
447
448
# File 'lib/raptor/http1.rb', line 446

def expects_100_continue?(env)
  (env[Rack::SERVER_PROTOCOL] == HTTP_11) && env[HTTP_EXPECT]&.casecmp?(EXPECT_100_CONTINUE)
end

#extract_body(buffer, env, parser, nread, decode_chunked:) ⇒ String, ...

Resolves the request body for a finished parse. Returns the body String (or nil when the request has no body), or one of :incomplete, :too_large, :malformed when the caller must fall back or reject.

RBS:

  • (String buffer, Hash[String, untyped] env, HttpParser parser, Integer nread, decode_chunked: bool) -> (String | Symbol)?

Parameters:

  • buffer (String)

    the raw request bytes

  • env (Hash)

    the Rack environment being built

  • parser (HttpParser)

    the parser holding the finished parse state

  • nread (Integer)

    the byte offset where the body begins in buffer

  • decode_chunked (Boolean)

    whether to decode chunked bodies inline; when false chunked bodies signal :incomplete

  • decode_chunked: (Boolean)

Returns:

  • (String, nil, Symbol)


416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
# File 'lib/raptor/http1.rb', line 416

def extract_body(buffer, env, parser, nread, decode_chunked:)
  return unless parser.has_body?

  body = buffer.byteslice(nread..-1) || ""

  if parser.chunked?
    return :incomplete unless decode_chunked

    body, chunked_state = HttpParser.decode_chunked(body, @max_body_size)
    case chunked_state
    when :complete
      env.delete(HTTP_TRANSFER_ENCODING)
      body
    else
      chunked_state
    end
  elsif parser.content_length > body.bytesize
    :incomplete
  else
    body
  end
end

#fallback_to_reactor(socket, id, buffer, env, parse_data, reactor, request_count, remote_addr, url_scheme, persisted: true) ⇒ void

This method returns an undefined value.

Re-registers a socket with the reactor for further processing when an incomplete request is received on the fast path.

RBS:

  • (TCPSocket socket, Integer id, String buffer, Hash[String, untyped] env, Hash[Symbol, untyped] parse_data, Reactor reactor, Integer request_count, String remote_addr, String url_scheme, persisted: bool) -> void

Parameters:

  • socket (TCPSocket)

    the client socket

  • id (Integer)

    unique client identifier

  • buffer (String)

    the partial request data already read

  • env (Hash)

    partial env hash from the HTTP parser

  • parse_data (Hash)

    metadata from the parsing pass

  • reactor (Reactor)

    the reactor to re-register with

  • request_count (Integer)

    number of requests handled on this connection

  • remote_addr (String)

    client IP address

  • url_scheme (String)

    "http" or "https"

  • persisted (Boolean) (defaults to: true)

    whether the connection has already completed at least one request

  • persisted: (Boolean) (defaults to: true)


676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
# File 'lib/raptor/http1.rb', line 676

def fallback_to_reactor(socket, id, buffer, env, parse_data, reactor, request_count, remote_addr, url_scheme, persisted: true)
  continued = expects_100_continue?(env)
  socket_write(socket, CONTINUE_RESPONSE) rescue nil if continued

  reactor.persist(socket, id, request_count, remote_addr: remote_addr, url_scheme: url_scheme)
  state = {
    id: id,
    buffer: buffer.dup,
    env: env,
    request_count: request_count,
    parse_data: parse_data,
    remote_addr: remote_addr,
    url_scheme: url_scheme
  }
  state[:persisted] = true if persisted
  state[:continued] = true if continued
  reactor.update_state(Ractor.make_shareable(state))
end

#forwarded_https?(env) ⇒ Boolean

Returns true when an upstream proxy signals that it terminated TLS for this request via X-Forwarded-Proto, X-Forwarded-Scheme, or X-Forwarded-Ssl. Only the first comma-separated value is consulted.

RBS:

  • (Hash[String, untyped] env) -> bool

Parameters:

  • env (Hash)

    the Rack environment

Returns:

  • (Boolean)


824
825
826
827
828
829
# File 'lib/raptor/http1.rb', line 824

def forwarded_https?(env)
  proto = env["HTTP_X_FORWARDED_PROTO"] || env["HTTP_X_FORWARDED_SCHEME"]
  return true if proto && proto.split(",").first&.strip&.casecmp?(Server::HTTPS_SCHEME)

  env["HTTP_X_FORWARDED_SSL"]&.casecmp?("on") || false
end

#handle_app_error(socket, rack_env, status, headers, error, response_started:, hijacked:) ⇒ void

This method returns an undefined value.

Handles an exception raised while processing a request. Fires the rack.response_finished callbacks with the error, writes a 500 response when no bytes have gone to the socket yet, and routes the exception through the configured on_error handler (or re-raises).

RBS:

  • (TCPSocket socket, Hash[String, untyped]? rack_env, Integer? status, Hash[String, String | Array[String]]? headers, Exception error, response_started: bool, hijacked: bool) -> void

Parameters:

  • socket (TCPSocket)

    the client socket

  • rack_env (Hash, nil)

    the Rack environment, if it was built

  • status (Integer, nil)

    the status returned by the app, if any

  • headers (Hash, nil)

    the headers returned by the app, if any

  • error (Exception)

    the exception raised

  • response_started (Boolean)

    whether any response bytes have been written

  • hijacked (Boolean)

    whether the app took over the socket

  • response_started: (Boolean)
  • hijacked: (Boolean)


562
563
564
565
566
567
568
569
570
571
# File 'lib/raptor/http1.rb', line 562

def handle_app_error(socket, rack_env, status, headers, error, response_started:, hijacked:)
  call_response_finished(rack_env, status, headers, error) if rack_env
  socket.write(INTERNAL_SERVER_ERROR_RESPONSE) rescue nil unless response_started || hijacked

  if @on_error
    @on_error.call(rack_env, error) rescue nil
  else
    raise error
  end
end

#handle_parsed_request(parsed_request, reactor, thread_pool) ⇒ void

This method returns an undefined value.

Dispatches a parsed HTTP request to the thread pool when complete, or hands it back to the reactor for more I/O when incomplete.

RBS:

  • (Hash[Symbol, untyped] parsed_request, Reactor reactor, AtomicThreadPool thread_pool) -> void

Parameters:

  • parsed_request (Hash)

    the parsed request state from the ractor pool

  • reactor (Reactor)

    the reactor managing the client connection

  • thread_pool (AtomicThreadPool)

    thread pool for application processing



325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
# File 'lib/raptor/http1.rb', line 325

def handle_parsed_request(parsed_request, reactor, thread_pool)
  if parsed_request[:too_large]
    socket = reactor.remove(parsed_request[:id])
    reject_oversized(socket) if socket
    return
  end

  if parsed_request[:malformed]
    socket = reactor.remove(parsed_request[:id])
    reject_malformed(socket) if socket
    return
  end

  unless parsed_request[:complete]
    parsed_request = send_continue_if_expected(parsed_request, reactor)
    reactor.update_state(parsed_request)
  else
    socket = reactor.remove(parsed_request[:id])
    request_count = (parsed_request[:request_count] || 0) + 1
    remote_addr = parsed_request[:remote_addr] || Server::DEFAULT_REMOTE_ADDR
    url_scheme = parsed_request[:url_scheme] || Server::HTTP_SCHEME

    thread_pool << proc do
      process_client(
        socket,
        parsed_request[:id],
        parsed_request[:env].dup,
        parsed_request[:parse_data],
        parsed_request[:body],
        reactor,
        thread_pool,
        request_count,
        remote_addr,
        url_scheme
      )
    end
  end
end

#keep_alive?(env, request_count) ⇒ Boolean

Returns true when the connection should be kept alive after the current response.

RBS:

  • (Hash[String, untyped] env, Integer request_count) -> bool

Parameters:

  • env (Hash)

    the Rack environment

  • request_count (Integer)

    number of requests handled on this connection

Returns:

  • (Boolean)

    true if the connection should be kept alive



839
840
841
842
843
844
845
846
847
848
849
# File 'lib/raptor/http1.rb', line 839

def keep_alive?(env, request_count)
  return false if request_count >= @max_keepalive_requests

  connection_header = env[HTTP_CONNECTION]

  if env[Rack::SERVER_PROTOCOL] == HTTP_11
    !connection_header&.casecmp?(CONNECTION_CLOSE)
  else
    connection_header&.casecmp?(CONNECTION_KEEPALIVE) || false
  end
end

#normalize_headers(headers) ⇒ Hash

Returns a normalised copy of the response headers with lowercased keys and illegal/rack.*/status entries dropped.

RBS:

  • (Hash[String, String | Array[String]] headers) -> Hash[String, String | Array[String]]

Parameters:

  • headers (Hash)

    raw headers from the Rack application

Returns:

  • (Hash)

    normalized headers with lowercased string keys

Raises:

  • (TypeError)

    if headers is not a Hash or a key is not a String



941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
# File 'lib/raptor/http1.rb', line 941

def normalize_headers(headers)
  raise TypeError, "headers must be a Hash" unless headers.is_a?(Hash)

  normalized = {}
  headers.each do |key, value|
    raise TypeError, "header keys must be Strings" unless key.is_a?(String)

    next if HttpParser.illegal_header_key?(key)

    normalized_key = key.match?(/[A-Z]/) ? key.downcase : key
    next if normalized_key.start_with?(RACK_HEADER_PREFIX)
    next if normalized_key == "status"

    normalized[normalized_key] = value
  end
  normalized
end

#parse_next_request(buffer) ⇒ Array(Hash, Hash, Integer, HttpParser)

Runs a fresh HTTP/1.x parse against buffer, returning [env, parse_data, nread, parser]. Raises HttpParserError on malformed input.

RBS:

  • (String buffer) -> [Hash[String, untyped], Hash[Symbol, untyped], Integer, HttpParser]

Parameters:

  • buffer (String)

    the raw request bytes

Returns:



394
395
396
397
398
399
400
401
# File 'lib/raptor/http1.rb', line 394

def parse_next_request(buffer)
  parser = (Thread.current[:raptor_http_parser] ||= HttpParser.new)
  parser.reset
  env = @env_template.dup
  nread = parser.execute(env, buffer, 0)
  parse_data = { parse_count: 1, content_length: parser.content_length }
  [env, parse_data, nread, parser]
end

#parser_workerProc

Instance-level wrapper around Raptor::Http.parser_worker that binds this handler's env template and body-size limit into the worker proc.

RBS:

  • () -> ^(Hash[Symbol, untyped]) -> Hash[Symbol, untyped]

Returns:

  • (Proc)


216
217
218
# File 'lib/raptor/http1.rb', line 216

def parser_worker
  Http.parser_worker(@env_template, @max_body_size)
end

#process_client(socket, id, env, parse_data, body, reactor, thread_pool, request_count, remote_addr, url_scheme) ⇒ void

This method returns an undefined value.

Processes a client connection by handling the current request and, if keep-alive, eagerly reading subsequent requests inline.

RBS:

  • (TCPSocket socket, Integer id, Hash[String, untyped] env, Hash[Symbol, untyped] parse_data, String? body, Reactor reactor, AtomicThreadPool thread_pool, Integer request_count, String remote_addr, String url_scheme) -> void

Parameters:

  • socket (TCPSocket)

    the client socket

  • id (Integer)

    unique client identifier

  • env (Hash)

    partial env hash from the HTTP parser

  • parse_data (Hash)

    metadata from the parsing pass

  • body (String, nil)

    decoded request body

  • reactor (Reactor)

    the reactor managing the client connection

  • thread_pool (AtomicThreadPool)

    thread pool for application processing

  • request_count (Integer)

    number of requests handled on this connection

  • remote_addr (String)

    client IP address

  • url_scheme (String)

    "http" or "https"



488
489
490
491
# File 'lib/raptor/http1.rb', line 488

def process_client(socket, id, env, parse_data, body, reactor, thread_pool, request_count, remote_addr, url_scheme)
  keep_alive = process_request(socket, env, parse_data, body, request_count, remote_addr, url_scheme)
  eager_keepalive(socket, id, reactor, thread_pool, request_count, remote_addr, url_scheme) if keep_alive
end

#process_request(socket, env, parse_data, body, request_count, remote_addr, url_scheme) ⇒ Boolean

Processes a single request. Builds the Rack env, calls the app, writes the response, and returns whether the connection stays open for another request.

RBS:

  • (TCPSocket socket, Hash[String, untyped] env, Hash[Symbol, untyped] parse_data, String? body, Integer request_count, String remote_addr, String url_scheme) -> bool

Parameters:

  • socket (TCPSocket)

    the client socket

  • env (Hash)

    partial env hash from the HTTP parser

  • parse_data (Hash)

    metadata from the parsing pass

  • body (String, nil)

    decoded request body

  • request_count (Integer)

    number of requests handled on this connection

  • remote_addr (String)

    client IP address

  • url_scheme (String)

    "http" or "https"

Returns:

  • (Boolean)

    true if the connection should be kept alive



507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
# File 'lib/raptor/http1.rb', line 507

def process_request(socket, env, parse_data, body, request_count, remote_addr, url_scheme)
  rack_env = nil
  status = nil
  headers = nil
  hijacked = false
  keep_alive = false
  response_started = false

  begin
    rack_env = build_rack_env(env, parse_data, body, socket, remote_addr: remote_addr, url_scheme: url_scheme)
    status, headers, body = @app.call(rack_env)

    if rack_env[RACK_HIJACKED]
      hijacked = true
      body.close if body.respond_to?(:close)
    else
      hijacked = headers.is_a?(Hash) && !!headers[Rack::RACK_HIJACK]
      streaming = body.respond_to?(:call) && !body.respond_to?(:each)
      keep_alive = (hijacked || streaming) ? false : keep_alive?(rack_env, request_count)
      response_size = response_size(headers, body) if @access_log_io && !hijacked
      response_started = true
      write_response(socket, rack_env, status, headers, body, keep_alive: keep_alive)
    end

    write_access_log(rack_env, status, response_size, remote_addr) if @access_log_io && !hijacked
    call_response_finished(rack_env, status, headers, nil)
    keep_alive && !hijacked
  rescue => error
    keep_alive = false
    handle_app_error(socket, rack_env, status, headers, error, response_started: response_started, hijacked: hijacked)
  ensure
    rack_input = rack_env && rack_env[Rack::RACK_INPUT]
    rack_input.close! rescue nil if rack_input.respond_to?(:close!)

    unless hijacked || keep_alive
      socket.close rescue nil
    end
  end
end

#read_into_thread_buffer(socket) ⇒ String

Reads pending bytes off socket into the thread-local read buffer, draining any additional SSL-buffered bytes so pending is empty on return. Raises IO::WaitReadable, EOFError, or IOError like the underlying read_nonblock does.

RBS:

  • (TCPSocket socket) -> String

Parameters:

  • socket (TCPSocket)

    the socket to read from

Returns:

  • (String)

    the thread-local buffer, freshly populated



375
376
377
378
379
380
381
382
383
384
# File 'lib/raptor/http1.rb', line 375

def read_into_thread_buffer(socket)
  buffer = (Thread.current[:raptor_read_buffer] ||= String.new(capacity: READ_BUFFER_SIZE))
  socket.read_nonblock(READ_BUFFER_SIZE, buffer)

  while socket.respond_to?(:pending) && socket.pending.positive?
    buffer << socket.read_nonblock(socket.pending)
  end

  buffer
end

#reject_malformed(socket) ⇒ void

This method returns an undefined value.

Writes a 400 response and closes the socket.

RBS:

  • (TCPSocket socket) -> void

Parameters:

  • socket (TCPSocket)

    the client socket



712
713
714
715
# File 'lib/raptor/http1.rb', line 712

def reject_malformed(socket)
  socket.write(BAD_REQUEST_RESPONSE) rescue nil
  socket.close rescue nil
end

#reject_oversized(socket) ⇒ void

This method returns an undefined value.

Writes a 413 response and closes the socket.

RBS:

  • (TCPSocket socket) -> void

Parameters:

  • socket (TCPSocket)

    the client socket



701
702
703
704
# File 'lib/raptor/http1.rb', line 701

def reject_oversized(socket)
  socket.write(CONTENT_TOO_LARGE_RESPONSE) rescue nil
  socket.close rescue nil
end

#response_size(headers, body) ⇒ String

Returns the response body size as a String for the access log, taken from the content-length header when set, computed from the body otherwise, or - when the size cannot be determined upfront.

RBS:

  • (Hash[String, String | Array[String]] headers, untyped body) -> String

Parameters:

  • headers (Hash)

    the response headers

  • body (Object)

    the response body

Returns:

  • (String)


1282
1283
1284
# File 'lib/raptor/http1.rb', line 1282

def response_size(headers, body)
  headers[Rack::CONTENT_LENGTH] || calculate_content_length(body)&.to_s || "-"
end

#send_continue_if_expected(state, reactor) ⇒ Hash

Sends an HTTP 100 Continue response when the client requested Expect: 100-continue, returning the state hash with :continued set once written. A write failure is silently ignored.

RBS:

  • (Hash[Symbol, untyped] state, Reactor reactor) -> Hash[Symbol, untyped]

Parameters:

  • state (Hash)

    the partially-parsed connection state

  • reactor (Reactor)

    the reactor holding the connection's socket

Returns:

  • (Hash)

    the state, with :continued set if 100 was written



459
460
461
462
463
464
465
466
467
468
469
470
# File 'lib/raptor/http1.rb', line 459

def send_continue_if_expected(state, reactor)
  return state if state[:continued]

  env = state[:env]
  return state unless env && expects_100_continue?(env)

  socket = reactor.socket_for(state[:id])
  return state unless socket

  socket_write(socket, CONTINUE_RESPONSE) rescue nil
  state.merge(continued: true)
end

#send_early_hints(socket, hints) ⇒ void

This method returns an undefined value.

Sends an HTTP 103 Early Hints response, skipping any entries with illegal header keys or values. No-ops when hints is empty.

RBS:

  • (TCPSocket socket, Hash[String, String | Array[String]] hints) -> void

Parameters:

  • socket (TCPSocket)

    the client socket to write to

  • hints (Hash)

    header name to value (or array of values) pairs



859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
# File 'lib/raptor/http1.rb', line 859

def send_early_hints(socket, hints)
  return if hints.empty?

  response = +"#{HTTP_11} 103 Early Hints\r\n"
  hints.each do |key, value|
    next if HttpParser.illegal_header_key?(key)

    values = value.is_a?(Array) ? value : [value]
    values.each do |hint_value|
      next if HttpParser.illegal_header_value?(hint_value.to_s)

      response << "#{key.downcase}: #{hint_value}\r\n"
    end
  end
  response << "\r\n"

  socket_write(socket, response)
end

#shutdownvoid

This method returns an undefined value.

Signals eager keep-alive loops to stop processing further requests on their connections. In-flight requests complete normally.

RBS:

  • () -> void



252
253
254
# File 'lib/raptor/http1.rb', line 252

def shutdown
  @running.make_false
end

#socket_write(socket, string) ⇒ void

This method returns an undefined value.

Instance-level wrapper around Raptor::Http.socket_write that applies the configured write_timeout.

RBS:

  • (TCPSocket socket, String string) -> void

Parameters:

  • socket (TCPSocket)

    the socket to write to

  • string (String)

    the data to write

Raises:

  • (Http::WriteError)

    if the socket is not writable within the timeout or raises IOError



229
230
231
# File 'lib/raptor/http1.rb', line 229

def socket_write(socket, string)
  Http.socket_write(socket, string, timeout: @write_timeout)
end

#socket_writev(socket, strings) ⇒ void

This method returns an undefined value.

Instance-level wrapper around Raptor::Http.socket_writev that applies the configured write_timeout.

RBS:

  • (TCPSocket socket, Array[String] strings) -> void

Parameters:

  • socket (TCPSocket)

    the socket to write to

  • strings (Array<String>)

    the buffers to write in order

Raises:

  • (Http::WriteError)

    if the socket is not writable within the timeout or raises IOError



242
243
244
# File 'lib/raptor/http1.rb', line 242

def socket_writev(socket, strings)
  Http.socket_writev(socket, strings, timeout: @write_timeout)
end

#uncork_socket(socket) ⇒ void

This method returns an undefined value.

Disables TCP_CORK on the socket, flushing any buffered packets. Linux-only; a no-op elsewhere.

RBS:

  • (TCPSocket socket) -> void

Parameters:

  • socket (TCPSocket)

    the socket to uncork



1305
1306
1307
# File 'lib/raptor/http1.rb', line 1305

def uncork_socket(socket)
  socket.setsockopt(Socket::IPPROTO_TCP, Socket::TCP_CORK, 0) if socket.is_a?(TCPSocket)
end

#validate_headers(headers, status, no_entity_body) ⇒ void

This method returns an undefined value.

Raises when the headers include entries forbidden for the response status (content-type or content-length on a 204, 304, or 1xx).

RBS:

  • (Hash[String, String | Array[String]] headers, Integer status, bool no_entity_body) -> void

Parameters:

  • headers (Hash)

    normalized response headers

  • status (Integer)

    HTTP status code

  • no_entity_body (Boolean)

    whether the status forbids an entity body

Raises:

  • (ArgumentError)

    if a forbidden header is present for the status



969
970
971
972
973
974
975
# File 'lib/raptor/http1.rb', line 969

def validate_headers(headers, status, no_entity_body)
  if no_entity_body
    raise ArgumentError, "content-type must not be present for status #{status}" if headers.key?(Rack::CONTENT_TYPE)

    raise ArgumentError, "content-length must not be present for status #{status}" if headers.key?(Rack::CONTENT_LENGTH)
  end
end

#validate_status(status) ⇒ void

This method returns an undefined value.

Validates that the status code is a valid integer.

RBS:

  • (Integer status) -> void

Parameters:

  • status (Object)

    the status value to validate

Raises:

  • (TypeError)

    if status is not an Integer

  • (ArgumentError)

    if status is less than 100



927
928
929
930
931
# File 'lib/raptor/http1.rb', line 927

def validate_status(status)
  raise TypeError, "status must be an Integer" unless status.is_a?(Integer)

  raise ArgumentError, "status must be >= 100" unless status >= 100
end

#write_access_log(env, status, size, remote_addr) ⇒ void

This method returns an undefined value.

Instance-level wrapper around Raptor::Http.write_access_log that routes to the configured @access_log_io.

RBS:

  • (Hash[String, untyped] env, Integer status, String size, String remote_addr) -> void

Parameters:

  • env (Hash)

    the Rack environment

  • status (Integer)

    the response status code

  • size (String)

    the response body size in bytes, or - if unknown

  • remote_addr (String)

    the client IP address



1269
1270
1271
# File 'lib/raptor/http1.rb', line 1269

def write_access_log(env, status, size, remote_addr)
  Http.write_access_log(@access_log_io, env, status, size, remote_addr)
end

#write_array_body(socket, response, body_array, use_chunked) ⇒ void

This method returns an undefined value.

Writes an array body to the socket.

RBS:

  • (TCPSocket socket, String response, Array[String] body_array, bool use_chunked) -> void

Parameters:

  • socket (TCPSocket)

    the client socket

  • response (String)

    headers already serialized, to be written before the body

  • body_array (Array<String>)

    the response body chunks

  • use_chunked (Boolean)

    whether to use chunked transfer encoding



1147
1148
1149
1150
1151
1152
1153
# File 'lib/raptor/http1.rb', line 1147

def write_array_body(socket, response, body_array, use_chunked)
  if body_array.length == 1
    write_single_chunk(socket, response, body_array.first, use_chunked)
  else
    write_multiple_chunks(socket, response, body_array, use_chunked)
  end
end

#write_enumerable_body(socket, response, body, use_chunked) ⇒ void

This method returns an undefined value.

Writes a generic enumerable body to the socket.

RBS:

  • (TCPSocket socket, String response, untyped body, bool use_chunked) -> void

Parameters:

  • socket (TCPSocket)

    the client socket

  • response (String)

    headers already serialized, to be written before the body

  • body (Object)

    any object responding to each

  • use_chunked (Boolean)

    whether to use chunked transfer encoding

Raises:

  • (TypeError)

    if any yielded chunk is not a String



1219
1220
1221
1222
1223
1224
1225
1226
1227
1228
1229
1230
1231
1232
1233
1234
1235
1236
1237
1238
1239
# File 'lib/raptor/http1.rb', line 1219

def write_enumerable_body(socket, response, body, use_chunked)
  if use_chunked
    socket_write(socket, response)
    buffer = +""
    body.each do |chunk|
      raise TypeError, "body must yield String values" unless chunk.is_a?(String)

      buffer.clear
      HttpParser.chunked_encode(buffer, chunk)
      socket_write(socket, buffer) unless buffer.empty?
    end
    socket_write(socket, CHUNKED_TERMINATOR)
  else
    body.each do |chunk|
      raise TypeError, "body must yield String values" unless chunk.is_a?(String)

      response << chunk
    end
    socket_write(socket, response)
  end
end

#write_file_body(socket, response, path, content_length, use_chunked) ⇒ void

This method returns an undefined value.

Writes a file body to the socket.

RBS:

  • (TCPSocket socket, String response, String path, Integer? content_length, bool use_chunked) -> void

Parameters:

  • socket (TCPSocket)

    the client socket

  • response (String)

    headers already serialized, to be written before the body

  • path (String)

    filesystem path of the file to send

  • content_length (Integer, nil)

    pre-calculated file size

  • use_chunked (Boolean)

    whether to use chunked transfer encoding



1115
1116
1117
1118
1119
1120
1121
1122
1123
1124
1125
1126
1127
1128
1129
1130
1131
1132
1133
1134
1135
1136
# File 'lib/raptor/http1.rb', line 1115

def write_file_body(socket, response, path, content_length, use_chunked)
  File.open(path, "rb") do |file|
    if use_chunked
      buffer = response
      while (chunk = file.read(FILE_CHUNK_SIZE))
        HttpParser.chunked_encode(buffer, chunk)
        if buffer.bytesize >= CHUNKED_WRITE_THRESHOLD
          socket_write(socket, buffer)
          buffer = +""
        end
      end
      buffer << CHUNKED_TERMINATOR
      socket_write(socket, buffer)
    elsif content_length && content_length < BODY_BUFFER_THRESHOLD
      response << file.read(content_length)
      socket_write(socket, response)
    else
      socket_write(socket, response)
      IO.copy_stream(file, socket)
    end
  end
end

#write_full_response(socket, response, headers, body, http_version) ⇒ void

This method returns an undefined value.

Writes a complete response with a body. Emits a Content-Length when the total size is known upfront, otherwise chunked encoding on HTTP/1.1.

RBS:

  • (TCPSocket socket, String response, Hash[String, String | Array[String]] headers, untyped body, String http_version) -> void

Parameters:

  • socket (TCPSocket)

    the client socket

  • response (String)

    the status line accumulated so far

  • headers (Hash)

    normalized response headers

  • body (Object)

    the response body

  • http_version (String)

    "HTTP/1.1" or "HTTP/1.0"



1043
1044
1045
1046
1047
1048
1049
1050
1051
1052
1053
1054
1055
1056
1057
1058
1059
1060
1061
1062
1063
1064
1065
1066
1067
1068
1069
1070
1071
1072
1073
1074
1075
1076
1077
1078
1079
1080
1081
1082
1083
# File 'lib/raptor/http1.rb', line 1043

def write_full_response(socket, response, headers, body, http_version)
  if body.respond_to?(:call)
    HttpParser.format_headers(response, headers)
    response << "\r\n"
    socket_write(socket, response)
    uncork_socket(socket)
    body.call(socket)
    return
  end

  content_length = headers[Rack::CONTENT_LENGTH]&.to_i
  use_chunked = false

  if !content_length || content_length.zero?
    calculated_length = calculate_content_length(body)
    if calculated_length
      content_length = calculated_length
    elsif http_version == HTTP_11 && !headers.key?(Rack::TRANSFER_ENCODING)
      use_chunked = true
    end
  end

  if content_length && content_length >= 0
    headers[Rack::CONTENT_LENGTH] = content_length.to_s
  elsif use_chunked
    headers[Rack::TRANSFER_ENCODING] = TRANSFER_ENCODING_CHUNKED
  end

  HttpParser.format_headers(response, headers)
  response << "\r\n"

  if body.respond_to?(:to_path) && (path = body.to_path) && File.readable?(path)
    write_file_body(socket, response, path, content_length, use_chunked)
  elsif body.respond_to?(:to_ary)
    write_array_body(socket, response, body.to_ary, use_chunked)
  elsif body.respond_to?(:each)
    write_enumerable_body(socket, response, body, use_chunked)
  else
    raise TypeError, "body must respond to each, to_ary, or to_path"
  end
end

#write_hijacked_response(socket, response, headers, response_hijack) ⇒ void

This method returns an undefined value.

Writes the response headers, uncorks the socket, and hands the raw socket to the hijack callback.

RBS:

  • (TCPSocket socket, String response, Hash[String, String | Array[String]] headers, ^(TCPSocket) -> void response_hijack) -> void

Parameters:

  • socket (TCPSocket)

    the client socket

  • response (String)

    the status line accumulated so far

  • headers (Hash)

    normalized response headers

  • response_hijack (Proc)

    callable that receives the socket and writes the body



1002
1003
1004
1005
1006
1007
1008
# File 'lib/raptor/http1.rb', line 1002

def write_hijacked_response(socket, response, headers, response_hijack)
  HttpParser.format_headers(response, headers)
  response << "\r\n"
  socket_write(socket, response)
  uncork_socket(socket)
  response_hijack.call(socket)
end

#write_multiple_chunks(socket, response, body_array, use_chunked) ⇒ void

This method returns an undefined value.

Writes a multi-element array body to the socket.

RBS:

  • (TCPSocket socket, String response, Array[String] body_array, bool use_chunked) -> void

Parameters:

  • socket (TCPSocket)

    the client socket

  • response (String)

    headers already serialized, to be written before the body

  • body_array (Array<String>)

    the response body chunks

  • use_chunked (Boolean)

    whether to use chunked transfer encoding

Raises:

  • (TypeError)

    if any chunk is not a String



1187
1188
1189
1190
1191
1192
1193
1194
1195
1196
1197
1198
1199
1200
1201
1202
1203
1204
1205
1206
1207
# File 'lib/raptor/http1.rb', line 1187

def write_multiple_chunks(socket, response, body_array, use_chunked)
  if use_chunked
    buffer = response
    body_array.each do |chunk|
      raise TypeError, "body must yield String values" unless chunk.is_a?(String)

      HttpParser.chunked_encode(buffer, chunk)
      if buffer.bytesize >= CHUNKED_WRITE_THRESHOLD
        socket_write(socket, buffer)
        buffer = +""
      end
    end
    buffer << CHUNKED_TERMINATOR
    socket_write(socket, buffer)
  else
    body_array.each do |chunk|
      raise TypeError, "body must yield String values" unless chunk.is_a?(String)
    end
    socket_writev(socket, [response, *body_array])
  end
end

#write_no_body_response(socket, response, headers, no_entity_body) ⇒ void

This method returns an undefined value.

Writes a response with no entity body, adding a zero Content-Length when the status may carry a body but none was supplied.

RBS:

  • (TCPSocket socket, String response, Hash[String, String | Array[String]] headers, bool no_entity_body) -> void

Parameters:

  • socket (TCPSocket)

    the client socket

  • response (String)

    the status line accumulated so far

  • headers (Hash)

    normalized response headers

  • no_entity_body (Boolean)

    whether the response status forbids an entity body



1021
1022
1023
1024
1025
1026
1027
1028
1029
# File 'lib/raptor/http1.rb', line 1021

def write_no_body_response(socket, response, headers, no_entity_body)
  unless no_entity_body
    headers[Rack::CONTENT_LENGTH] = "0" unless headers.key?(Rack::CONTENT_LENGTH) || headers.key?(Rack::TRANSFER_ENCODING)
  end

  HttpParser.format_headers(response, headers)
  response << "\r\n"
  socket_write(socket, response)
end

#write_response(socket, env, status, headers, body, keep_alive: false) ⇒ void

This method returns an undefined value.

Writes a complete HTTP response to the socket.

RBS:

  • (TCPSocket socket, Hash[String, untyped] env, Integer status, Hash[String, String | Array[String]] headers, untyped body, ?keep_alive: bool) -> void

Parameters:

  • socket (TCPSocket)

    the client socket to write to

  • env (Hash)

    the Rack environment

  • status (Integer)

    HTTP status code

  • headers (Hash)

    response headers from the Rack application

  • body (Object)

    response body (array, enumerable, file, or callable)

  • keep_alive (Boolean) (defaults to: false)

    whether to send a keep-alive connection header

  • keep_alive: (Boolean) (defaults to: false)


889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
# File 'lib/raptor/http1.rb', line 889

def write_response(socket, env, status, headers, body, keep_alive: false)
  validate_status(status)
  response_hijack = headers.is_a?(Hash) ? headers.delete(Rack::RACK_HIJACK) : nil
  headers = normalize_headers(headers)
  no_entity_body = (status >= 100 && status < 200) || status == 204 || status == 304
  validate_headers(headers, status, no_entity_body)

  headers["connection"] = keep_alive ? CONNECTION_KEEPALIVE : CONNECTION_CLOSE

  http_version = env[Rack::SERVER_PROTOCOL] == HTTP_11 ? HTTP_11 : HTTP_10
  no_body = env[Rack::REQUEST_METHOD] == "HEAD" || no_entity_body

  response = build_status_line(http_version, status)

  corked = !keep_alive
  cork_socket(socket) if corked

  if response_hijack
    write_hijacked_response(socket, response, headers, response_hijack)
  elsif no_body
    write_no_body_response(socket, response, headers, no_entity_body)
  else
    write_full_response(socket, response, headers, body, http_version)
  end
ensure
  body.close if body.respond_to?(:close)
  uncork_socket(socket) if corked
  socket.flush rescue nil
end

#write_single_chunk(socket, response, chunk, use_chunked) ⇒ void

This method returns an undefined value.

Writes a single-element array body to the socket.

RBS:

  • (TCPSocket socket, String response, String chunk, bool use_chunked) -> void

Parameters:

  • socket (TCPSocket)

    the client socket

  • response (String)

    headers already serialized, to be written before the body

  • chunk (String)

    the single body chunk

  • use_chunked (Boolean)

    whether to use chunked transfer encoding

Raises:

  • (TypeError)

    if the chunk is not a String



1165
1166
1167
1168
1169
1170
1171
1172
1173
1174
1175
# File 'lib/raptor/http1.rb', line 1165

def write_single_chunk(socket, response, chunk, use_chunked)
  raise TypeError, "body must yield String values" unless chunk.is_a?(String)

  if use_chunked
    HttpParser.chunked_encode(response, chunk)
    response << CHUNKED_TERMINATOR
    socket_write(socket, response)
  else
    socket_writev(socket, [response, chunk])
  end
end