Class: HTTPX::Connection::HTTP2
- Inherits:
-
Object
- Object
- HTTPX::Connection::HTTP2
show all
- Includes:
- HTTPX::Callbacks, Loggable
- Defined in:
- lib/httpx/connection/http2.rb,
sig/connection/http2.rbs
Defined Under Namespace
Classes: Error, GoawayError, PingError
Constant Summary
collapse
- MAX_CONCURRENT_REQUESTS =
::HTTP2::DEFAULT_MAX_CONCURRENT_STREAMS
Constants included
from Loggable
Loggable::COLORS, Loggable::USE_DEBUG_LOG
Instance Attribute Summary collapse
Instance Method Summary
collapse
-
#<<(data) ⇒ void
-
#can_buffer_more_requests? ⇒ Boolean
-
#close ⇒ void
-
#consume ⇒ void
-
#emit_stream_error(stream, request, error) ⇒ void
-
#empty? ⇒ Boolean
-
#end_stream?(request, next_chunk) ⇒ Boolean
-
#exhausted? ⇒ Boolean
-
#frame_with_extra_info(frame) ⇒ Hash[Symbol, untyped]
-
#handle(request, stream) ⇒ void
-
#handle_error(ex, request = nil) ⇒ void
-
#handle_stream(stream, request) ⇒ void
-
#init_connection ⇒ void
(also: #reset)
-
#initialize(buffer, options) ⇒ HTTP2
constructor
-
#interests ⇒ io_interests?
-
#join_body(stream, request) ⇒ void
-
#join_headers(stream, request) ⇒ void
-
#join_headline ⇒ String
-
#join_trailers(stream, request) ⇒ void
-
#on_altsvc(origin, frame) ⇒ void
-
#on_close(_last_frame, error, _payload) ⇒ void
-
#on_frame(bytes) ⇒ void
-
#on_frame_received(frame) ⇒ void
-
#on_frame_sent(frame) ⇒ void
-
#on_origin(origin) ⇒ void
-
#on_pong(ping) ⇒ void
-
#on_promise(stream) ⇒ void
-
#on_settings ⇒ void
-
#on_stream_close(stream, request, error) ⇒ void
-
#on_stream_data(stream, request, data) ⇒ void
-
#on_stream_half_close(stream, request) ⇒ void
-
#on_stream_headers(stream, request, h) ⇒ void
-
#on_stream_refuse(stream, request, error) ⇒ void
-
#on_stream_trailers(stream, request, response, h) ⇒ void
-
#ping ⇒ void
-
#reset_requests ⇒ void
-
#send(request, head = false) ⇒ Boolean
-
#send_chunk(request, stream, chunk, next_chunk) ⇒ void
-
#send_pending ⇒ void
-
#set_protocol_headers(request) ⇒ _Each[[String, String]]
-
#teardown(request = nil) ⇒ void
-
#timeout ⇒ interval?
-
#waiting_for_ping? ⇒ Boolean
Methods included from Loggable
#log, #log_exception, log_identifiers, #log_redact, #log_redact_body, #log_redact_headers
#callbacks, #callbacks_for?, #emit, #on, #once
Constructor Details
#initialize(buffer, options) ⇒ HTTP2
Returns a new instance of HTTP2.
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
|
# File 'lib/httpx/connection/http2.rb', line 33
def initialize(buffer, options)
@callbacks = nil
@options = options
@settings = @options.http2_settings
@pending = []
@streams = {}
@drains = {}
@pings = []
@streams_to_close_after_receive = []
@buffer = buffer
@handshake_completed = false
@wait_for_handshake = @settings.key?(:wait_for_handshake) ? @settings.delete(:wait_for_handshake) : true
@max_concurrent_requests = @options.max_concurrent_requests || MAX_CONCURRENT_REQUESTS
@max_requests = @options.max_requests
init_connection
end
|
Instance Attribute Details
#pending ⇒ Array[Request]
Returns the value of attribute pending.
31
32
33
|
# File 'lib/httpx/connection/http2.rb', line 31
def pending
@pending
end
|
#streams ⇒ Hash[Request, ::HTTP2::Stream]
Returns the value of attribute streams.
31
32
33
|
# File 'lib/httpx/connection/http2.rb', line 31
def streams
@streams
end
|
Instance Method Details
#<<(data) ⇒ void
This method returns an undefined value.
113
114
115
116
117
118
119
120
121
|
# File 'lib/httpx/connection/http2.rb', line 113
def <<(data)
@connection << data
while (stream, request, error = @streams_to_close_after_receive.shift)
emit_stream_error(stream, request, error)
end
end
|
#can_buffer_more_requests? ⇒ Boolean
186
187
188
189
190
|
# File 'lib/httpx/connection/http2.rb', line 186
def can_buffer_more_requests?
(@handshake_completed || !@wait_for_handshake) &&
@streams.size < @max_concurrent_requests &&
@streams.size < @max_requests
end
|
#close ⇒ void
This method returns an undefined value.
97
98
99
100
101
102
103
|
# File 'lib/httpx/connection/http2.rb', line 97
def close
unless @connection.closed?
@connection.goaway
emit(:timeout, @options.timeout[:close_handshake_timeout])
end
emit(:close)
end
|
#consume ⇒ void
This method returns an undefined value.
143
144
145
146
147
148
149
|
# File 'lib/httpx/connection/http2.rb', line 143
def consume
@streams.each do |request, stream|
next unless request.can_buffer?
handle(request, stream)
end
end
|
#emit_stream_error(stream, request, error) ⇒ void
This method returns an undefined value.
534
535
536
537
538
|
# File 'lib/httpx/connection/http2.rb', line 534
def emit_stream_error(stream, request, error)
teardown(request)
stream.close
emit(:error, request, error)
end
|
#empty? ⇒ Boolean
105
106
107
|
# File 'lib/httpx/connection/http2.rb', line 105
def empty?
@connection.closed? || @streams.empty?
end
|
#end_stream?(request, next_chunk) ⇒ Boolean
304
305
306
|
# File 'lib/httpx/connection/http2.rb', line 304
def end_stream?(request, next_chunk)
!(next_chunk || request.trailers? || request.callbacks_for?(:trailers))
end
|
#exhausted? ⇒ Boolean
109
110
111
|
# File 'lib/httpx/connection/http2.rb', line 109
def exhausted?
!@max_requests.positive?
end
|
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
|
# File 'lib/httpx/connection/http2.rb', line 468
def (frame)
flags_bits = frame.fetch(:flags, 0)
case frame[:type]
when :data
flags = [] flags << :end_stream if flags_bits.anybits?(0b0001)
flags << :padded if flags_bits.anybits?(0b1000)
frame.merge(payload: frame[:payload].bytesize, flags: flags)
when :push_promise, :headers
flags = [] flags << :end_stream if flags_bits.anybits?(0b0001)
flags << :priority if flags_bits.anybits?(0b0010)
flags << :end_headers if flags_bits.anybits?(0b0100)
flags << :padded if flags_bits.anybits?(0b1000)
frame.merge(payload: (frame[:payload]), flags: flags)
when :ping
flags = [] flags << :ack if flags_bits.anybits?(0b0001)
frame.merge(payload: (frame[:payload]), flags: flags)
when :settings
flags = [] flags << :ack if flags_bits.anybits?(0b0001)
frame.merge(flags: flags)
when :window_update
connection_or_stream = if (id = frame[:stream]).zero?
@connection
else
@streams.each_value.find { |s| s.id == id }
end
if connection_or_stream
frame.merge(
local_window: connection_or_stream.local_window,
remote_window: connection_or_stream.remote_window,
buffered_amount: connection_or_stream.buffered_amount,
stream_state: connection_or_stream.state,
)
else
frame
end
else
frame
end.merge(connection_state: @connection.state)
end
|
#handle(request, stream) ⇒ void
This method returns an undefined value.
198
199
200
201
202
203
204
205
206
207
208
|
# File 'lib/httpx/connection/http2.rb', line 198
def handle(request, stream)
catch(:buffer_full) do
request.transition(:headers)
(stream, request) if request.state == :headers
request.transition(:body)
join_body(stream, request) if request.state == :body
request.transition(:trailers)
join_trailers(stream, request) if request.state == :trailers && !request.body.empty?
request.transition(:done)
end
end
|
#handle_error(ex, request = nil) ⇒ void
This method returns an undefined value.
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
|
# File 'lib/httpx/connection/http2.rb', line 151
def handle_error(ex, request = nil)
if ex.is_a?(OperationTimeoutError) && !@handshake_completed && @connection.state != :closed
@connection.goaway(:settings_timeout, "closing due to settings timeout")
emit(:close_handshake)
settings_ex = SettingsTimeoutError.new(ex.timeout, ex.message)
settings_ex.set_backtrace(ex.backtrace)
ex = settings_ex
end
while (req, _ = @streams.shift)
next if request && request == req
emit(:error, req, ex)
end
while (req = @pending.shift)
next if request && request == req
emit(:error, req, ex)
end
end
|
#handle_stream(stream, request) ⇒ void
This method returns an undefined value.
233
234
235
236
237
238
239
240
|
# File 'lib/httpx/connection/http2.rb', line 233
def handle_stream(stream, request)
request.on(:refuse, &method(:on_stream_refuse).curry(3)[stream, request])
stream.on(:close, &method(:on_stream_close).curry(3)[stream, request])
stream.on(:half_close) { on_stream_half_close(stream, request) }
stream.on(:altsvc, &method(:on_altsvc).curry(2)[request.origin])
stream.on(:headers, &method(:on_stream_headers).curry(3)[stream, request])
stream.on(:data, &method(:on_stream_data).curry(3)[stream, request])
end
|
#init_connection ⇒ void
Also known as:
reset
This method returns an undefined value.
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
|
# File 'lib/httpx/connection/http2.rb', line 210
def init_connection
@connection = ::HTTP2::Client.new(@settings)
@connection.on(:frame, &method(:on_frame))
@connection.on(:frame_sent, &method(:on_frame_sent))
@connection.on(:frame_received, &method(:on_frame_received))
@connection.on(:origin, &method(:on_origin))
@connection.on(:promise, &method(:on_promise))
@connection.on(:altsvc) { |frame| on_altsvc(frame[:origin], frame) }
@connection.on(:settings_ack, &method(:on_settings))
@connection.on(:ack, &method(:on_pong))
@connection.on(:goaway, &method(:on_close))
@connection.send_connection_preface
end
|
#interests ⇒ io_interests?
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
|
# File 'lib/httpx/connection/http2.rb', line 56
def interests
if @connection.closed?
return unless @handshake_completed
return if @buffer.empty?
return :w
end
unless @connection.state == :connected && @handshake_completed
return @buffer.empty? ? :r : :rw
end
unless @connection.send_buffer.empty?
return :rw unless @buffer.empty?
return :r
end
return :w if !@pending.empty? && can_buffer_more_requests?
return :w unless @drains.empty?
if @buffer.empty?
return if @streams.empty? && @pings.empty?
:r
else
:w
end
end
|
#join_body(stream, request) ⇒ void
This method returns an undefined value.
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
|
# File 'lib/httpx/connection/http2.rb', line 277
def join_body(stream, request)
return if request.body.empty?
chunk = @drains.delete(request) || request.drain_body
while chunk
next_chunk = request.drain_body
send_chunk(request, stream, chunk, next_chunk)
if next_chunk && (@buffer.full? || request.body.unbounded_body?)
@drains[request] = next_chunk
throw(:buffer_full)
end
chunk = next_chunk
end
return unless (error = request.drain_error)
on_stream_refuse(stream, request, error)
end
|
This method returns an undefined value.
251
252
253
254
255
256
257
258
259
260
261
262
263
|
# File 'lib/httpx/connection/http2.rb', line 251
def (stream, request)
= (request)
if request..key?("host")
request.log { "forbidden \"host\" header found (#{(request.["host"])}), will use it as authority..." }
[":authority"] = request.["host"]
end
request.log(level: 1, color: :yellow) do
"\n#{request..merge().each.map { |k, v| "#{stream.id}: -> HEADER: #{k}: #{(v)}" }.join("\n")}"
end
stream.(request..each(), end_stream: request.body.empty?)
end
|
#join_headline ⇒ String
66
|
# File 'sig/connection/http2.rbs', line 66
def join_headline: (Request request) -> String
|
#join_trailers(stream, request) ⇒ void
This method returns an undefined value.
265
266
267
268
269
270
271
272
273
274
275
|
# File 'lib/httpx/connection/http2.rb', line 265
def join_trailers(stream, request)
unless request.trailers?
stream.data("", end_stream: true) if request.callbacks_for?(:trailers)
return
end
request.log(level: 1, color: :yellow) do
request.trailers.each.map { |k, v| "#{stream.id}: -> HEADER: #{k}: #{(v)}" }.join("\n")
end
stream.(request.trailers.each, end_stream: true)
end
|
#on_altsvc(origin, frame) ⇒ void
This method returns an undefined value.
512
513
514
515
516
517
518
|
# File 'lib/httpx/connection/http2.rb', line 512
def on_altsvc(origin, frame)
log(level: 2) { "#{frame[:stream]}: altsvc frame was received" }
log(level: 2) { "#{frame[:stream]}: #{(frame.inspect)}" }
alt_origin = URI.parse("#{frame[:proto]}://#{frame[:host]}:#{frame[:port]}")
params = { "ma" => frame[:max_age] }
emit(:altsvc, origin, alt_origin, origin, params)
end
|
#on_close(_last_frame, error, _payload) ⇒ void
This method returns an undefined value.
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
|
# File 'lib/httpx/connection/http2.rb', line 435
def on_close(_last_frame, error, _payload)
is_connection_closed = @connection.closed?
if error
@buffer.clear if is_connection_closed
case error
when :http_1_1_required
while (request = @pending.shift)
emit(:error, request, error)
end
else
ex = GoawayError.new(error)
ex.set_backtrace(caller)
handle_error(ex)
teardown
end
end
return unless is_connection_closed && @streams.empty?
emit(:close) if is_connection_closed
end
|
#on_frame(bytes) ⇒ void
This method returns an undefined value.
424
425
426
|
# File 'lib/httpx/connection/http2.rb', line 424
def on_frame(bytes)
@buffer << bytes
end
|
#on_frame_received(frame) ⇒ void
This method returns an undefined value.
463
464
465
466
|
# File 'lib/httpx/connection/http2.rb', line 463
def on_frame_received(frame)
log(level: 2) { "#{frame[:stream]}: frame was received!" }
log(level: 2, color: :magenta) { "#{frame[:stream]}: #{(frame)}" }
end
|
#on_frame_sent(frame) ⇒ void
This method returns an undefined value.
458
459
460
461
|
# File 'lib/httpx/connection/http2.rb', line 458
def on_frame_sent(frame)
log(level: 2) { "#{frame[:stream]}: frame was sent!" }
log(level: 2, color: :blue) { "#{frame[:stream]}: #{(frame)}" }
end
|
#on_origin(origin) ⇒ void
This method returns an undefined value.
524
525
526
|
# File 'lib/httpx/connection/http2.rb', line 524
def on_origin(origin)
emit(:origin, origin)
end
|
#on_pong(ping) ⇒ void
This method returns an undefined value.
528
529
530
531
532
|
# File 'lib/httpx/connection/http2.rb', line 528
def on_pong(ping)
raise PingError unless @pings.delete(ping.to_s)
emit(:pong)
end
|
#on_promise(stream) ⇒ void
This method returns an undefined value.
520
521
522
|
# File 'lib/httpx/connection/http2.rb', line 520
def on_promise(stream)
emit(:promise, @streams.key(stream.parent), stream)
end
|
#on_settings ⇒ void
This method returns an undefined value.
428
429
430
431
432
433
|
# File 'lib/httpx/connection/http2.rb', line 428
def on_settings(*)
@handshake_completed = true
emit(:current_timeout)
@max_concurrent_requests = [@max_concurrent_requests, @connection.remote_settings[:settings_max_concurrent_streams]].min
send_pending
end
|
#on_stream_close(stream, request, error) ⇒ void
This method returns an undefined value.
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
|
# File 'lib/httpx/connection/http2.rb', line 387
def on_stream_close(stream, request, error)
return if error == :stream_closed && !@streams.key?(request)
log(level: 2) { "#{stream.id}: closing stream" }
teardown(request)
if error
case error
when :http_1_1_required
emit(:error, request, error)
else
ex = Error.new(stream.id, error)
ex.set_backtrace(caller)
response = ErrorResponse.new(request, ex)
request.response = response
emit(:response, request, response)
end
else
if (response = request.response)
if response.is_a?(Response) && response.status == 421
emit(:error, request, :http_1_1_required)
else
emit(:response, request, response)
end
end
end
send(@pending.shift) unless @pending.empty?
return unless @streams.empty? && exhausted?
if @pending.empty?
close
else
emit(:exhausted)
end
end
|
#on_stream_data(stream, request, data) ⇒ void
This method returns an undefined value.
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
|
# File 'lib/httpx/connection/http2.rb', line 355
def on_stream_data(stream, request, data)
request.log(level: 1, color: :green) { "#{stream.id}: <- DATA: #{data.bytesize} bytes..." }
request.log(level: 2, color: :green) { "#{stream.id}: <- #{log_redact_body(data.inspect)}" }
return unless request.response
request.response << data
rescue HTTPX::Error => e
if stream.state == :closing
emit_stream_error(stream, request, e)
else
@streams_to_close_after_receive << [stream, request, e]
end
end
|
#on_stream_half_close(stream, request) ⇒ void
This method returns an undefined value.
377
378
379
380
381
382
383
384
385
|
# File 'lib/httpx/connection/http2.rb', line 377
def on_stream_half_close(stream, request)
unless stream.send_buffer.empty?
stream.send_buffer.clear
stream.data("", end_stream: true)
end
request.log(level: 2) { "#{stream.id}: waiting for response..." }
end
|
This method returns an undefined value.
HTTP/2 Callbacks
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
|
# File 'lib/httpx/connection/http2.rb', line 312
def (stream, request, h)
response = request.response
if response.is_a?(Response) && response.version == "2.0"
on_stream_trailers(stream, request, response, h)
return
end
request.log(color: :yellow) do
h.map { |k, v| "#{stream.id}: <- HEADER: #{k}: #{k == ":status" ? v : (v)}" }.join("\n")
end
_, status = h.shift
= request.options..new(h)
raise HTTPX::Error, "maximum number of response headers exceeded" if h.size > @options.
if ( = @options.)
.each do |_, v| raise HTTPX::Error, "maximum header value size exceeded" if v.size >
end
end
response = request.options.response_class.new(request, status, "2.0", )
if response.content_length && response.content_length > request.options.max_response_body_size
raise HTTPX::Error.new, "maximum response body size exceeded"
end
request.response = response
@streams[request] = stream
handle(request, stream) if request.expects?
rescue HTTPX::Error => e
@streams_to_close_after_receive << [stream, request, e]
end
|
#on_stream_refuse(stream, request, error) ⇒ void
This method returns an undefined value.
372
373
374
375
|
# File 'lib/httpx/connection/http2.rb', line 372
def on_stream_refuse(stream, request, error)
on_stream_close(stream, request, error)
stream.close
end
|
#on_stream_trailers(stream, request, response, h) ⇒ void
This method returns an undefined value.
348
349
350
351
352
353
|
# File 'lib/httpx/connection/http2.rb', line 348
def on_stream_trailers(stream, request, response, h)
request.log(color: :yellow) do
h.map { |k, v| "#{stream.id}: <- HEADER: #{k}: #{(v)}" }.join("\n")
end
response.(h)
end
|
#ping ⇒ void
This method returns an undefined value.
171
172
173
174
175
176
|
# File 'lib/httpx/connection/http2.rb', line 171
def ping
ping = SecureRandom.gen_random(8)
@connection.ping(ping.dup)
ensure
@pings << ping
end
|
#reset_requests ⇒ void
This method returns an undefined value.
182
|
# File 'lib/httpx/connection/http2.rb', line 182
def reset_requests; end
|
#send(request, head = false) ⇒ Boolean
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
|
# File 'lib/httpx/connection/http2.rb', line 123
def send(request, head = false)
unless can_buffer_more_requests?
head ? @pending.unshift(request) : @pending << request
return false
end
unless (stream = @streams[request])
stream = @connection.new_stream(**request.http2_stream_options)
handle_stream(stream, request)
@streams[request] = stream
@max_requests -= 1
end
handle(request, stream)
true
rescue ::HTTP2::Error::StreamLimitExceeded
@pending.unshift(request)
false
rescue ::HTTP2::Error::Error, ArgumentError => e
emit(:error, request, e)
end
|
#send_chunk(request, stream, chunk, next_chunk) ⇒ void
This method returns an undefined value.
298
299
300
301
302
|
# File 'lib/httpx/connection/http2.rb', line 298
def send_chunk(request, stream, chunk, next_chunk)
request.log(level: 1, color: :green) { "#{stream.id}: -> DATA: #{chunk.bytesize} bytes..." }
request.log(level: 2, color: :green) { "#{stream.id}: -> #{log_redact_body(chunk.inspect)}" }
stream.data(chunk, end_stream: end_stream?(request, next_chunk))
end
|
#send_pending ⇒ void
This method returns an undefined value.
192
193
194
195
196
|
# File 'lib/httpx/connection/http2.rb', line 192
def send_pending
while (request = @pending.shift)
break unless send(request, true)
end
end
|
242
243
244
245
246
247
248
249
|
# File 'lib/httpx/connection/http2.rb', line 242
def (request)
{
":scheme" => request.scheme,
":method" => request.verb,
":path" => request.path,
":authority" => request.authority,
}
end
|
#teardown(request = nil) ⇒ void
This method returns an undefined value.
540
541
542
543
544
545
546
547
548
|
# File 'lib/httpx/connection/http2.rb', line 540
def teardown(request = nil)
if request
@drains.delete(request)
@streams.delete(request)
else
@drains.clear
@streams.clear
end
end
|
#timeout ⇒ interval?
50
51
52
53
54
|
# File 'lib/httpx/connection/http2.rb', line 50
def timeout
return @options.timeout[:operation_timeout] if @handshake_completed
@options.timeout[:settings_timeout]
end
|
#waiting_for_ping? ⇒ Boolean
178
179
180
|
# File 'lib/httpx/connection/http2.rb', line 178
def waiting_for_ping?
@pings.any?
end
|