Class: Raptor::Reactor

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

Overview

Multiplexes client connections and closes them when they overrun their per-phase timeouts (first data, subsequent chunks, and keep-alive idle). Feeds ready bytes to the parser pipeline and hands completed requests back to the caller-provided handlers.

Defined Under Namespace

Classes: TimeoutClient

Constant Summary collapse

CHUNK_SIZE =

Returns:

  • (Object)
64 * 1024
TIMEOUT_RESPONSE =

Returns:

  • (::String)
"HTTP/1.1 408 Request Timeout\r\nContent-Length: 0\r\nConnection: close\r\n\r\n"

Instance Method Summary collapse

Constructor Details

#initialize(ractor_pool, thread_pool, connection_options:, http1_options:) ⇒ Reactor

Creates a new Reactor instance.

Parameters:

  • ractor_pool (RactorPool)

    ractor pool for HTTP parsing

  • thread_pool (AtomicThreadPool)

    thread pool for application processing

  • connection_options (Hash)

    per-connection timeout configuration

  • http1_options (Hash)

    HTTP/1.1-specific configuration

  • connection_options: (Hash[Symbol, untyped])
  • http1_options: (Hash[Symbol, untyped])

Options Hash (connection_options:):

  • :first_data_timeout (Integer)

    timeout for initial data

  • :chunk_data_timeout (Integer)

    timeout for subsequent chunks

Options Hash (http1_options:):

  • :persistent_data_timeout (Integer)

    timeout for keep-alive idle connections



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

def initialize(ractor_pool, thread_pool, connection_options:, http1_options:)
  @ractor_pool = ractor_pool
  @thread_pool = thread_pool
  @first_data_timeout = connection_options[:first_data_timeout]
  @chunk_data_timeout = connection_options[:chunk_data_timeout]
  @persistent_data_timeout = http1_options[:persistent_data_timeout]

  @selector = NIO::Selector.new
  @queue = Queue.new
  @timeouts = RedBlackTree.new

  @id_to_socket = {}
  @socket_to_state = {}
  @id_to_timeout = {}
  @id_to_writer = {}
  @id_to_flow_control = {}
end

Instance Method Details

#add(state) ⇒ void

This method returns an undefined value.

Adds a new client connection to the reactor.

Parameters:

  • state (Hash)

    client connection state including socket and ID

Options Hash (state):

  • :socket (TCPSocket)

    the client socket

  • :id (Integer)

    unique identifier for the client



157
158
159
160
161
162
163
164
165
166
167
168
# File 'lib/raptor/reactor.rb', line 157

def add(state)
  socket = state[:socket]
  state.delete(:socket)
  writer = state.delete(:writer)
  flow_control = state.delete(:flow_control)
  @id_to_socket[state[:id]] = socket
  @socket_to_state[socket] = state
  @id_to_writer[state[:id]] = writer if writer
  @id_to_flow_control[state[:id]] = flow_control if flow_control

  read_and_queue_for_parse(socket, state)
end

#backlogInteger

Returns the number of complete requests either being processed or awaiting processing.

Returns:

  • (Integer)

    number of complete requests



315
316
317
# File 'lib/raptor/reactor.rb', line 315

def backlog
  @thread_pool.queue_size + @thread_pool.active_count
end

#cleanup(socket) ⇒ void

This method returns an undefined value.

Cleans up a client connection by removing it from tracking and closing the socket.

Parameters:

  • socket (TCPSocket)

    the socket to clean up



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

def cleanup(socket)
  state = @socket_to_state.delete(socket)
  @id_to_socket.delete(state[:id])
  @id_to_writer.delete(state[:id])
  @id_to_flow_control.delete(state[:id])
  socket.close
end

#close_connection(id) ⇒ void

This method returns an undefined value.

Closes the socket for the given connection and drops all reactor state associated with it.

Parameters:

  • id (Integer)

    unique client identifier



288
289
290
291
292
293
294
295
296
# File 'lib/raptor/reactor.rb', line 288

def close_connection(id)
  socket = @id_to_socket.delete(id)
  return unless socket

  @socket_to_state.delete(socket)
  @id_to_writer.delete(id)
  @id_to_flow_control.delete(id)
  socket.close rescue nil
end

#complete?(state) ⇒ Boolean

Returns true when the request has been fully parsed.

Parameters:

  • state (Hash)

    connection state

Returns:

  • (Boolean)

    true if the request is complete



409
410
411
# File 'lib/raptor/reactor.rb', line 409

def complete?(state)
  state[:complete]
end

#first_data_received?(state) ⇒ Boolean

Checks if any data has been received for this connection.

Parameters:

  • state (Hash)

    connection state

Returns:

  • (Boolean)

    true if first data has been received



419
420
421
# File 'lib/raptor/reactor.rb', line 419

def first_data_received?(state)
  complete?(state) || (state.dig(:parse_data, :parse_count) || 0) >= 1
end

#flow_control_for(id) ⇒ Object?

Returns the flow controller associated with a given connection, if one was supplied when the connection was added.

Parameters:

  • id (Integer)

    unique client identifier

Returns:

  • (Object, nil)

    the flow controller, if found



259
260
261
# File 'lib/raptor/reactor.rb', line 259

def flow_control_for(id)
  @id_to_flow_control[id]
end

#persist(socket, id, request_count, remote_addr:, url_scheme:) ⇒ void

This method returns an undefined value.

Re-registers a kept-alive connection for the next request cycle under the persistent-data timeout.

Parameters:

  • socket (TCPSocket)

    the kept-alive client socket

  • id (Integer)

    the unique client identifier

  • request_count (Integer)

    number of requests handled on this connection

  • remote_addr (String)

    the client's remote IP address

  • url_scheme (String)

    "http" or "https"

  • remote_addr: (String)
  • url_scheme: (String)


214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
# File 'lib/raptor/reactor.rb', line 214

def persist(socket, id, request_count, remote_addr:, url_scheme:)
  state = {
    id: id,
    request_count: request_count,
    remote_addr: remote_addr,
    url_scheme: url_scheme,
    persisted: true
  }

  @id_to_socket[id] = socket
  @socket_to_state[socket] = state
  @queue << socket
  @selector.wakeup
rescue ClosedQueueError
  socket.close
end

#read_and_queue_for_parse(socket, state) ⇒ Hash?

Reads data from a socket and either queues it for parsing or returns it to the selector.

Parameters:

  • socket (TCPSocket)

    the socket to read from and queue

  • state (Hash)

    current connection state

Returns:

  • (Hash, nil)

    updated state, if successful



366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
# File 'lib/raptor/reactor.rb', line 366

def read_and_queue_for_parse(socket, state)
  data = begin
    socket.read_nonblock(CHUNK_SIZE)
  rescue IO::WaitReadable
    @queue << socket
    @selector.wakeup
    return
  rescue EOFError
    cleanup(socket)
    return
  end

  buffer = state[:buffer] ? state[:buffer].dup : String.new
  buffer << data

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

  state = state.frozen? ? state.merge(buffer: buffer) : state.merge!(buffer: buffer)
  @ractor_pool << Ractor.make_shareable(state)
end

#register(socket) ⇒ void

This method returns an undefined value.

Registers a socket with the NIO selector and sets up timeout tracking.

Parameters:

  • socket (TCPSocket)

    the socket to register



327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
# File 'lib/raptor/reactor.rb', line 327

def register(socket)
  @selector.register(socket, :r).value = socket

  state = @socket_to_state[socket]
  client = TimeoutClient.new(state)
  timeout = if state[:persisted]
    @persistent_data_timeout
  elsif first_data_received?(state)
    @chunk_data_timeout
  else
    @first_data_timeout
  end
  client.timeout_at = Process.clock_gettime(Process::CLOCK_MONOTONIC) + timeout
  @timeouts << client
  @id_to_timeout[state[:id]] = client
end

#remove(id) ⇒ TCPSocket?

Drops the reactor's references to a client whose parsed request has been handed off to the thread pool. The socket itself is kept open so the worker can write the response.

Parameters:

  • id (Integer)

    unique client identifier

Returns:

  • (TCPSocket, nil)

    the socket associated with id, if any



197
198
199
200
201
# File 'lib/raptor/reactor.rb', line 197

def remove(id)
  @id_to_socket.delete(id).tap do |socket|
    @socket_to_state.delete(socket)
  end
end

#runThread

Starts the reactor's main event loop in a new thread. Runs until the registration queue is closed and drained.

Returns:

  • (Thread)

    the thread running the reactor event loop



106
107
108
109
110
111
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
# File 'lib/raptor/reactor.rb', line 106

def run
  Thread.new do
    Thread.current.name = "Reactor"

    until @queue.closed? && @queue.empty?
      begin
        timeout = @timeouts.min&.timeout(Process.clock_gettime(Process::CLOCK_MONOTONIC))
        @selector.select(timeout) do |monitor|
          wakeup!(monitor.value)
        end

        now = Process.clock_gettime(Process::CLOCK_MONOTONIC)
        expired = []
        @timeouts.traverse do |to_client|
          break unless to_client.timeout(now) == 0

          expired << to_client
        end

        expired.each do |to_client|
          @timeouts.delete!(to_client)
          id = to_client.client_data[:id]
          @id_to_timeout.delete(id)
          socket = @id_to_socket[id]
          next unless socket

          @selector.deregister(socket)
          socket.write(TIMEOUT_RESPONSE) rescue nil
          cleanup(socket)
        end

        until @queue.empty?
          register(@queue.pop)
        end
      rescue => error
        Log.rescued_error(error)
      end
    end

    @selector.close
  end
end

#shutdownvoid

This method returns an undefined value.

Closes the registration queue and wakes the selector so the event loop drains pending work and exits.



304
305
306
307
# File 'lib/raptor/reactor.rb', line 304

def shutdown
  @queue.close
  @selector.wakeup
end

#socket_for(id) ⇒ TCPSocket?

Returns the socket for a given client identifier without removing it.

Parameters:

  • id (Integer)

    unique client identifier

Returns:

  • (TCPSocket, nil)

    the socket, if found



237
238
239
# File 'lib/raptor/reactor.rb', line 237

def socket_for(id)
  @id_to_socket[id]
end

#update_http2_state(state) ⇒ void

This method returns an undefined value.

Updates connection state for an HTTP/2 connection and re-registers the socket for further reads.

Parameters:

  • state (Hash)

    updated connection state from the ractor pool



270
271
272
273
274
275
276
277
278
279
# File 'lib/raptor/reactor.rb', line 270

def update_http2_state(state)
  socket = @id_to_socket[state[:id]]
  return unless socket

  @socket_to_state[socket] = state
  @queue << socket
  @selector.wakeup
rescue ClosedQueueError
  socket.close
end

#update_state(state) ⇒ void

This method returns an undefined value.

Updates the state of an existing client connection and re-registers it for further I/O.

Parameters:

  • state (Hash)

    updated client connection state

Options Hash (state):

  • :id (Integer)

    client identifier



178
179
180
181
182
183
184
185
186
187
# File 'lib/raptor/reactor.rb', line 178

def update_state(state)
  socket = @id_to_socket[state[:id]]
  return unless socket

  @socket_to_state[socket] = state
  @queue << socket
  @selector.wakeup
rescue ClosedQueueError
  socket.close
end

#wakeup!(socket) ⇒ void

This method returns an undefined value.

Handles socket wakeup by deregistering and queuing for processing.

Parameters:

  • socket (TCPSocket)

    the socket that became ready



350
351
352
353
354
355
356
# File 'lib/raptor/reactor.rb', line 350

def wakeup!(socket)
  @selector.deregister(socket)
  state = @socket_to_state[socket]
  to_client = @id_to_timeout.delete(state[:id])
  @timeouts.delete!(to_client)
  read_and_queue_for_parse(socket, state)
end

#writer_for(id) ⇒ Object?

Returns the writer object associated with a given connection, if one was supplied when the connection was added.

Parameters:

  • id (Integer)

    unique client identifier

Returns:

  • (Object, nil)

    the writer, if found



248
249
250
# File 'lib/raptor/reactor.rb', line 248

def writer_for(id)
  @id_to_writer[id]
end