Methods
Public Class
Public Instance
Classes and Modules
Constants
| MAX_CONCURRENT_REQUESTS | = | ::HTTP2::DEFAULT_MAX_CONCURRENT_STREAMS |
Public Instance Aliases
| reset | -> | init_connection |
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
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