class HTTPX::Connection::HTTP2

  1. lib/httpx/connection/http2.rb
Superclass: Object

Included modules

  1. Callbacks
  2. Loggable

Constants

MAX_CONCURRENT_REQUESTS = ::HTTP2::DEFAULT_MAX_CONCURRENT_STREAMS  

Public Instance Aliases

reset -> init_connection

Attributes

pending [R]
streams [R]

Public Class methods

new(buffer, options)
[show source]
   # File lib/httpx/connection/http2.rb
53 def initialize(buffer, options)
54   @callbacks = nil
55   @options = options
56   @settings = @options.http2_settings
57   @pending = []
58   @streams = {}
59   @drains = {}
60   @pings = []
61   @streams_to_close_after_receive = []
62   @buffer = buffer
63   @handshake_completed = false
64   @wait_for_handshake = @settings.key?(:wait_for_handshake) ? @settings.delete(:wait_for_handshake) : true
65   @max_concurrent_requests = @options.max_concurrent_requests || MAX_CONCURRENT_REQUESTS
66   @max_requests = @options.max_requests
67   init_connection
68 end

Public Instance methods

<<(data)
[show source]
    # File lib/httpx/connection/http2.rb
141 def <<(data)
142   @connection << data
143 
144   while (stream, request, error = @streams_to_close_after_receive.shift)
145     # these streams were marked for cancellation due to errors found while processing the
146     # data received by the peer.
147     emit_stream_error(stream, request, error)
148   end
149 end
close()
[show source]
    # File lib/httpx/connection/http2.rb
125 def close
126   unless @connection.closed?
127     @connection.goaway
128     emit(:timeout, @options.timeout[:close_handshake_timeout])
129   end
130   emit(:close)
131 end
consume()
[show source]
    # File lib/httpx/connection/http2.rb
171 def consume
172   @streams.each do |request, stream|
173     next unless request.can_buffer?
174 
175     handle(request, stream)
176   end
177 end
empty?()
[show source]
    # File lib/httpx/connection/http2.rb
133 def empty?
134   @connection.closed? || @streams.empty?
135 end
exhausted?()
[show source]
    # File lib/httpx/connection/http2.rb
137 def exhausted?
138   !@max_requests.positive?
139 end
handle_error(ex, request = nil)
[show source]
    # File lib/httpx/connection/http2.rb
179 def handle_error(ex, request = nil)
180   last_stream_id = Float::INFINITY
181   case ex
182   when OperationTimeoutError
183     if !@handshake_completed && @connection.state != :closed
184       @connection.goaway(:settings_timeout, "closing due to settings timeout")
185       emit(:close_handshake)
186       settings_ex = SettingsTimeoutError.new(ex.timeout, ex.message)
187       settings_ex.set_backtrace(ex.backtrace)
188       ex = settings_ex
189     end
190   when GoawayError
191     last_stream_id = ex.last_stream_id
192   end
193 
194   inflight_unprocessed_requests = [] #: Array[Request]
195 
196   while (req, stream = @streams.shift)
197     next if request && request == req
198 
199     if stream.id > last_stream_id
200       req.transition(:idle)
201       # unprocessed request
202       inflight_unprocessed_requests << req
203 
204       next
205     end
206 
207     emit(:error, req, ex)
208   end
209 
210   if ex.is_a?(GoawayError) && !ex.unrecoverable?
211     # resend unprocessed requests on a different connection
212     @pending.unshift(*inflight_unprocessed_requests) if inflight_unprocessed_requests.any?
213     emit(:exhausted, ex) if @pending.any?
214     return
215   end
216 
217   while (req = @pending.shift)
218     next if request && request == req
219 
220     emit(:error, req, ex)
221   end
222 end
interests()
[show source]
    # File lib/httpx/connection/http2.rb
 76 def interests
 77   if @connection.closed?
 78     return unless @handshake_completed
 79 
 80     return if @buffer.empty?
 81 
 82     # HTTP/2 GOAWAY frame buffered.
 83     return :w
 84   end
 85 
 86   unless @connection.state == :connected && @handshake_completed
 87     # HTTP/2 in intermediate state or still completing initialization-
 88     return @buffer.empty? ? :r : :rw
 89   end
 90 
 91   unless @connection.send_buffer.empty?
 92     # HTTP/2 connection is buffering data chunks and failing to emit DATA frames,
 93     # most likely because the flow control window is exhausted.
 94     return :rw unless @buffer.empty?
 95 
 96     # waiting for WINDOW_UPDATE frames
 97     return :r
 98   end
 99 
100   # only wait for writable if we're not waiting on more streams,
101   # or waiting on ping ACKs.
102   nothing_to_wait_for = @streams.empty? && @pings.empty?
103 
104   # there are pending bufferable requests
105   if !@pending.empty? && can_buffer_more_requests?
106     # only wait for writable if we're not waiting on more streams,
107     # or waiting on ping ACKs.
108     return nothing_to_wait_for ? :w : :rw
109   end
110 
111   # there are pending frames from the last run
112   return :w unless @drains.empty?
113 
114   if @buffer.empty?
115     # skip if no more requests or pings to process
116     return if nothing_to_wait_for
117 
118     :r
119   else
120     # buffered frames
121     nothing_to_wait_for ? :w : :rw
122   end
123 end
ping()
[show source]
    # File lib/httpx/connection/http2.rb
224 def ping
225   ping = SecureRandom.gen_random(8)
226   @connection.ping(ping.dup)
227 ensure
228   @pings << ping
229 end
reset_requests()
[show source]
    # File lib/httpx/connection/http2.rb
235 def reset_requests; end
send(request, head = false)
[show source]
    # File lib/httpx/connection/http2.rb
151 def send(request, head = false)
152   unless can_buffer_more_requests?
153     head ? @pending.unshift(request) : @pending << request
154     return false
155   end
156   unless (stream = @streams[request])
157     stream = @connection.new_stream(**request.http2_stream_options)
158     handle_stream(stream, request)
159     @streams[request] = stream
160     @max_requests -= 1
161   end
162   handle(request, stream)
163   true
164 rescue ::HTTP2::Error::StreamLimitExceeded
165   @pending.unshift(request)
166   false
167 rescue ::HTTP2::Error::Error, ArgumentError => e
168   emit(:error, request, e)
169 end
timeout()
[show source]
   # File lib/httpx/connection/http2.rb
70 def timeout
71   return @options.timeout[:operation_timeout] if @handshake_completed
72 
73   @options.timeout[:settings_timeout]
74 end
waiting_for_ping?()
[show source]
    # File lib/httpx/connection/http2.rb
231 def waiting_for_ping?
232   @pings.any?
233 end