From d7b80a6d2664770bd621ab986508c4a5f4210351 Mon Sep 17 00:00:00 2001 From: ilyazub Date: Fri, 2 Oct 2026 22:23:28 +0200 Subject: [PATCH] Resend idempotent requests when a reused connection was closed A persistent client keeps its socket between requests, and the server or a proxy may close it while it sits idle. The kernel still accepts the next request into the send buffer of the half-closed socket, so the failure only surfaces when the response is read: "couldn't read response headers" on a plain socket, ECONNRESET after a reset, and OpenSSL::SSL::SSLError over TLS on OpenSSL 3. In that case the server never saw the request (#420, #459). RFC 9110 Section 9.2.2 lets a client repeat an idempotent request automatically, and Net::HTTP, Go's net/http and urllib3 do so by default. When a reused connection fails before any response byte arrives and the request is replayable, close the connection and send the request once more on a new one. The resend starts from a fresh connection, so it happens at most once. Request#replayable? is true when the method is idempotent or an Idempotency-Key or X-Idempotency-Key header is present, and the body is nil or a String; IO and Enumerable bodies may already be consumed. Clients with a retriable policy keep its semantics. Features see one on_request and one wrap_response per call, and on_error only for the final failure. The client stays dirty until the resent request completes, so a resend interrupted by Thread#kill or Timeout is not reused. To stay within the Metrics limits, transmit and resend? live in Client::ConnectionReuse, check_premature_eof moves into Connection::Internals, the idempotency rules move into Request::Idempotency, and notify_features is inlined. --- CHANGELOG.md | 17 + lib/http/client.rb | 20 +- lib/http/client/connection_reuse.rb | 45 ++ lib/http/connection.rb | 29 +- lib/http/connection/internals.rb | 30 +- lib/http/request.rb | 2 + lib/http/request/idempotency.rb | 31 ++ sig/http.rbs | 31 +- test/http/client_resend_test.rb | 391 ++++++++++++++++++ test/http/connection_test.rb | 58 +++ test/http/request_test.rb | 44 ++ .../connection_reuse_tests.rb | 37 +- test/support/scripted_server.rb | 156 +++++++ 13 files changed, 850 insertions(+), 41 deletions(-) create mode 100644 lib/http/client/connection_reuse.rb create mode 100644 lib/http/request/idempotency.rb create mode 100644 test/http/client_resend_test.rb create mode 100644 test/support/scripted_server.rb diff --git a/CHANGELOG.md b/CHANGELOG.md index 954aa1ae..dbd2f0ec 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,8 +7,23 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +### Added + +- `HTTP::Request#replayable?` reports whether a request can be sent again + after a connection failure: its method is idempotent (RFC 9110 Section 9.2.2) + or it carries an `Idempotency-Key` / `X-Idempotency-Key` header, and its body + is nil or a String. + ### Fixed +- Persistent clients now resend a replayable request once, on a new + connection, when a reused connection fails before any response byte arrives, + typically because the server or a proxy closed it while it sat idle. + Previously the request raised `HTTP::ResponseHeaderError` ("couldn't read + response headers"), `HTTP::SocketReadError`, `HTTP::SocketWriteError` or, + over TLS on OpenSSL 3, `OpenSSL::SSL::SSLError`. Non-replayable requests, + clients with a `retriable` policy, timeouts, and failures after part of the + response arrived still raise. ([#420], [#459]) - Building a default `Host` header now raises `HTTP::RequestError` when the request URI has a nil host (previously `NoMethodError`) or an empty host (e.g. `https:///path` or `https://:123/path`, which previously produced @@ -298,9 +313,11 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 [#358]: https://github.com/httprb/http/issues/358 [#371]: https://github.com/httprb/http/issues/371 [#372]: https://github.com/httprb/http/issues/372 +[#420]: https://github.com/httprb/http/issues/420 [#447]: https://github.com/httprb/http/issues/447 [#448]: https://github.com/httprb/http/issues/448 [#449]: https://github.com/httprb/http/issues/449 +[#459]: https://github.com/httprb/http/issues/459 [#491]: https://github.com/httprb/http/issues/491 [#493]: https://github.com/httprb/http/pull/493 [#512]: https://github.com/httprb/http/issues/512 diff --git a/lib/http/client.rb b/lib/http/client.rb index d000be7d..abc4e209 100644 --- a/lib/http/client.rb +++ b/lib/http/client.rb @@ -1,12 +1,14 @@ # frozen_string_literal: true require "forwardable" +require "openssl" require "http/form_data" require "http/retriable/performer" require "http/options" require "http/feature" require "http/headers" +require "http/client/connection_reuse" require "http/connection" require "http/redirector" require "http/request/builder" @@ -17,6 +19,7 @@ module HTTP class Client extend Forwardable include Chainable + include ConnectionReuse # Initialize a new HTTP Client # @@ -142,14 +145,8 @@ def perform_with_retry(req, options) # @return [void] # @api private def send_request(req, options) - notify_features(req, options) - - @connection ||= HTTP::Connection.new(req, options) - - unless @connection.failed_proxy_connect? - @connection.send_request(req) - @connection.read_headers! - end + options.features.each_value { |feature| feature.on_request(req) } + transmit(req, options) rescue Error => e options.features.each_value { |feature| feature.on_error(req, e) } raise @@ -166,13 +163,6 @@ def build_wrapped_response(req, options) end end - # Notify features of an upcoming request attempt - # @return [void] - # @api private - def notify_features(req, options) - options.features.each_value { |feature| feature.on_request(req) } - end - # Execute the HTTP exchange wrapped by feature around_request hooks # @return [HTTP::Response] the response # @api private diff --git a/lib/http/client/connection_reuse.rb b/lib/http/client/connection_reuse.rb new file mode 100644 index 00000000..69e1091d --- /dev/null +++ b/lib/http/client/connection_reuse.rb @@ -0,0 +1,45 @@ +# frozen_string_literal: true + +module HTTP + class Client + # Sends requests over the client's connection, once more on a new + # connection when a reused one turns out to be closed + module ConnectionReuse + private + + # Write the request and read the response headers + # + # A reused connection may have been closed by the peer while it sat + # idle. When that surfaces before any response byte arrives, a + # replayable request is sent once more on a fresh connection. The resend + # always starts from a new connection, so it happens at most once. + # + # @return [void] + # @api private + def transmit(req, options) + reused = !@connection.nil? + @connection ||= Connection.new(req, options) + return if @connection.failed_proxy_connect? + + @connection.send_request(req) + @connection.read_headers! + rescue ConnectionError, OpenSSL::SSL::SSLError + raise unless reused && resend?(req, options) + + @connection.close + @connection = nil + retry + end + + # Whether a request that failed on a reused connection can be resent + # + # Explicit retry policies set with {Chainable#retriable} take precedence. + # + # @return [Boolean] + # @api private + def resend?(req, options) + !options.retriable && req.replayable? && !@connection.response_started? + end + end + end +end diff --git a/lib/http/connection.rb b/lib/http/connection.rb index e550a4b6..89bc3ed9 100644 --- a/lib/http/connection.rb +++ b/lib/http/connection.rb @@ -107,7 +107,8 @@ def send_request(req) raise StateError, "Tried to send a request while a response is pending. Make sure you read off the body." end - @pending_request = true + @response_started = false + @pending_request = true req.stream @socket @@ -199,6 +200,17 @@ def finished_request? !@pending_request && !@pending_response end + # Whether response bytes have arrived since the last request was sent + # + # @example + # connection.response_started? # => false + # + # @return [Boolean] + # @api public + def response_started? + @response_started + end + # Whether we're keeping the conn alive # # @example @@ -231,25 +243,12 @@ def init_state(options) @keep_alive_timeout = options.keep_alive_timeout.to_f @pending_request = false @pending_response = false + @response_started = false @failed_proxy_connect = false @buffer = "".b @parser = Response::Parser.new end - # Check for premature end-of-file and raise if detected - # - # @example - # check_premature_eof(:eof) - # - # @return [void] - # @api private - def check_premature_eof(eof) - return unless eof && !@parser.finished? && body_framed? - - close - raise ConnectionError, "response body ended prematurely" - end - # Connect socket and set up proxy/TLS # @return [void] # @api private diff --git a/lib/http/connection/internals.rb b/lib/http/connection/internals.rb index 3cb662eb..1bf3ed27 100644 --- a/lib/http/connection/internals.rb +++ b/lib/http/connection/internals.rb @@ -107,6 +107,20 @@ def set_keep_alive end end + # Check for premature end-of-file and raise if detected + # + # @example + # check_premature_eof(:eof) + # + # @return [void] + # @api private + def check_premature_eof(eof) + return unless eof && !@parser.finished? && body_framed? + + close + raise ConnectionError, "response body ended prematurely" + end + # Check if the response body has a known framing mechanism # # @example @@ -131,11 +145,25 @@ def read_more(size) @parser << "" :eof elsif value - @parser << value + feed_parser(value) end rescue IOError, SocketError, SystemCallError => e raise SocketReadError, "error reading from socket: #{e}", e.backtrace end + + # Marks the response as started, then feeds a chunk into parser + # + # The mark comes first so a chunk the parser rejects still counts. + # + # @example + # feed_parser("HTTP/1.1 200 OK\r\n") + # + # @return [Response::Parser] + # @api private + def feed_parser(chunk) + @response_started = true + @parser << chunk + end end end end diff --git a/lib/http/request.rb b/lib/http/request.rb index 31bad432..f782c0db 100644 --- a/lib/http/request.rb +++ b/lib/http/request.rb @@ -7,6 +7,7 @@ require "http/errors" require "http/headers" require "http/request/body" +require "http/request/idempotency" require "http/request/proxy" require "http/request/writer" require "http/version" @@ -19,6 +20,7 @@ class Request include HTTP::Base64 include Proxy + include Idempotency # The method given was not understood class UnsupportedMethodError < RequestError; end diff --git a/lib/http/request/idempotency.rb b/lib/http/request/idempotency.rb new file mode 100644 index 00000000..2bc07712 --- /dev/null +++ b/lib/http/request/idempotency.rb @@ -0,0 +1,31 @@ +# frozen_string_literal: true + +module HTTP + class Request + # Decides whether a request can be sent again on a new connection + module Idempotency + # Idempotent methods (RFC 9110, Section 9.2.2) + IDEMPOTENT_METHODS = %i[get head options trace put delete].freeze + + # Headers that mark any request as idempotent (draft-ietf-httpapi-idempotency-key-header) + IDEMPOTENCY_KEY_HEADERS = %w[Idempotency-Key X-Idempotency-Key].freeze + + # Whether the request can be sent again after a connection failure + # + # True when the method is idempotent or an idempotency key header is + # present, and the body is nil or a String, so it can be written again. + # + # @example + # request.replayable? # => true + # + # @return [Boolean] + # @api public + def replayable? + source = body.source + return false unless source.nil? || source.is_a?(String) + + IDEMPOTENT_METHODS.include?(verb) || IDEMPOTENCY_KEY_HEADERS.any? { |name| headers.include?(name) } + end + end + end +end diff --git a/sig/http.rbs b/sig/http.rbs index 9b1c4a69..88310c8d 100644 --- a/sig/http.rbs +++ b/sig/http.rbs @@ -127,6 +127,7 @@ module HTTP class Client extend Forwardable include Chainable + include Client::ConnectionReuse @default_options: Options @connection: untyped @@ -166,13 +167,23 @@ module HTTP def perform_once: (Request req, Options options) -> Response def perform_with_retry: (Request req, Options options) -> Response - def notify_features: (Request req, Options options) -> void def perform_exchange: (Request req, Options options) -> Response def around_request: (Request request, Options options) { (Request) -> Response } -> Response def build_response: (Request req, Options options) -> Response def build_wrapped_response: (Request req, Options options) -> Response def send_request: (Request req, Options options) -> void + def transmit: (Request req, Options options) -> void + def resend?: (Request req, Options options) -> bool def verify_connection!: (URI uri) -> void + + module ConnectionReuse + @connection: untyped + + private + + def transmit: (Request req, Options options) -> void + def resend?: (Request req, Options options) -> bool + end end class Session @@ -805,6 +816,7 @@ module HTTP extend Forwardable include Base64 include Request::Proxy + include Request::Idempotency class Builder HTTP_OR_HTTPS_RE: Regexp @@ -891,6 +903,19 @@ module HTTP def port: () -> untyped end + module Idempotency : _IdempotencyHost + interface _IdempotencyHost + def body: () -> Request::Body + def verb: () -> verb + def headers: () -> Headers + end + + IDEMPOTENT_METHODS: Array[verb] + IDEMPOTENCY_KEY_HEADERS: Array[String] + + def replayable?: () -> bool + end + class Body @source: untyped @@ -1127,6 +1152,7 @@ module HTTP @keep_alive_timeout: Float @pending_request: bool @pending_response: bool | Response + @response_started: bool @failed_proxy_connect: bool @buffer: String @parser: Response::Parser @@ -1145,6 +1171,7 @@ module HTTP def finish_response: () -> void def close: () -> void def finished_request?: () -> bool + def response_started?: () -> bool def keep_alive?: () -> bool def expired?: () -> bool def status_code: () -> Integer? @@ -1179,8 +1206,10 @@ module HTTP def handle_proxy_connect_response: () -> void def reset_timer: () -> void def set_keep_alive: () -> void + def check_premature_eof: (bool eof) -> void def body_framed?: () -> bool def read_more: (Integer size) -> untyped + def feed_parser: (String chunk) -> Response::Parser end end diff --git a/test/http/client_resend_test.rb b/test/http/client_resend_test.rb new file mode 100644 index 00000000..d649569d --- /dev/null +++ b/test/http/client_resend_test.rb @@ -0,0 +1,391 @@ +# frozen_string_literal: true + +require "test_helper" + +require "stringio" + +require "support/scripted_server" + +class HTTPClientResendTest < Minitest::Test + cover "HTTP::Client*" + + CountingFeature = Class.new(HTTP::Feature) do + attr_reader :requests, :responses, :errors + + def initialize + super + @requests = [] + @responses = 0 + @errors = 0 + end + + def on_request(request) + @requests << request.uri.path + end + + def wrap_response(response) + @responses += 1 + response + end + + def on_error(_request, _error) + @errors += 1 + end + end + + # Unwinds like Thread#exit: ensure clauses run, rescue clauses do not + InterruptingFeature = Class.new(HTTP::Feature) do + def wrap_response(response) + throw :interrupted if response.request.uri.path == "/interrupt" + + response + end + end + + def teardown + @client&.close + @server&.shutdown + @proxy&.shutdown + super + end + + def test_resends_get_after_idle_connection_is_closed + start(idle_close(:fin), serve) + build_client + get + @server.wait_closed + + assert_equal "ok", get + assert_equal 2, @server.accepts + end + + def test_resends_get_after_idle_connection_is_reset + start(idle_close(:reset), serve) + build_client + get + @server.wait_closed + + assert_equal "ok", get + assert_equal 2, @server.accepts + end + + def test_resends_get_after_idle_close_with_unread_body + start(idle_close(:fin), serve) + build_client + @client.get("#{@server.endpoint}/") + @server.wait_closed + + assert_equal "ok", get + assert_equal 2, @server.accepts + end + + def test_resends_get_after_unread_body_too_big_to_flush + size = HTTP::Connection::MAX_FLUSH_SIZE + 1 + start(answer("HTTP/1.1 200 OK\r\nContent-Length: #{size}\r\n\r\n#{'x' * size}"), serve) + build_client + @client.get("#{@server.endpoint}/") + + assert_equal "ok", get + assert_equal 2, @server.accepts + end + + def test_resends_identical_get_when_connection_closes_after_reading_it + start(close_after_second_request, serve) + build_client + get("/first") + + assert_equal "ok", get("/second") + + _, sent, resent = @server.requests + + assert_match %r{\AGET /second }, sent + assert_equal sent, resent + assert_equal 2, @server.accepts + end + + def test_resends_post_with_idempotency_key + start(close_after_second_request, serve) + build_client + get + response = @client.post("#{@server.endpoint}/", headers: { "Idempotency-Key" => "abc" }, body: "data") + _, sent, resent = @server.requests + + assert_equal "ok", response.to_s + assert_match(/^Idempotency-Key: abc\r$/i, sent) + assert sent.end_with?("\r\n\r\ndata") + assert_equal sent, resent + assert_equal 3, @server.requests.size + assert_equal 2, @server.accepts + end + + def test_does_not_resend_post + start(close_after_second_request, serve) + build_client + get + + assert_raises(HTTP::ConnectionError) { @client.post("#{@server.endpoint}/", body: "data") } + assert_equal 2, @server.requests.size + assert_equal 1, @server.accepts + end + + def test_does_not_resend_request_with_io_body + start(close_after_second_request, serve) + build_client + get + + assert_raises(HTTP::ConnectionError) { @client.put("#{@server.endpoint}/", body: StringIO.new("data")) } + assert_equal 2, @server.requests.size + assert_equal 1, @server.accepts + end + + def test_does_not_resend_after_partial_response + start(reply_to_second_request("HTTP/1.1 200 OK\r\nContent-"), serve) + build_client + get + + assert_raises(HTTP::ConnectionError) { get } + assert_equal 1, @server.accepts + end + + def test_does_not_resend_after_malformed_response + start(reply_to_second_request("garbage\r\n\r\n"), serve) + build_client + get + + assert_raises(HTTP::ConnectionError) { get } + assert_equal 1, @server.accepts + end + + def test_does_not_resend_after_informational_response + start(reply_to_second_request("HTTP/1.1 100 Continue\r\n\r\n"), serve) + build_client + get + + assert_raises(HTTP::ConnectionError) { get } + assert_equal 1, @server.accepts + end + + def test_does_not_resend_after_read_timeout + start(on_second_request(&:read_request), serve) + build_client(timeout_class: HTTP::Timeout::PerOperation, timeout_options: { read_timeout: 0.2 }) + get + + assert_raises(HTTP::TimeoutError) { get } + assert_equal 1, @server.accepts + end + + def test_does_not_resend_on_fresh_connection + start(read_then_close, serve) + build_client + + assert_raises(HTTP::ConnectionError) { get } + assert_equal 1, @server.accepts + end + + def test_resends_at_most_once + start(close_after_second_request, read_then_close) + build_client + get + + assert_raises(HTTP::ConnectionError) { get } + assert_equal 3, @server.requests.size + assert_equal 2, @server.accepts + end + + def test_closes_dead_connection_before_resending + start(close_after_second_request, serve) + build_client + get + dead = @client.instance_variable_get(:@connection).instance_variable_get(:@socket) + + assert_equal "ok", get + assert_predicate dead, :closed? + end + + def test_does_not_reuse_connection_after_resend_is_interrupted + start(close_after_second_request, serve) + build_client(features: { interrupting: InterruptingFeature.new }) + get + catch(:interrupted) { get("/interrupt") } + + assert_equal "ok", get + assert_equal 3, @server.accepts + end + + def test_raises_when_reconnecting_fails + start(close_after_second_request, serve) + build_client + get + + HTTP::Connection.stub(:new, ->(*) { raise HTTP::ConnectionError, "refused" }) do + err = assert_raises(HTTP::ConnectionError) { get } + + assert_equal "refused", err.message + end + + assert_equal "ok", get + assert_equal 2, @server.accepts + end + + def test_does_not_resend_when_retriable + start(close_after_second_request, serve) + build_client(retriable: { tries: 1 }) + get + + assert_raises(HTTP::OutOfRetriesError) { get } + assert_equal 1, @server.accepts + end + + def test_retriable_retries_on_a_new_connection + start(close_after_second_request, serve) + build_client(retriable: { tries: 2, delay: 0 }) + get + + assert_equal "ok", get + assert_equal 2, @server.accepts + end + + def test_features_see_each_call_once_when_resend_succeeds + feature = CountingFeature.new + start(close_after_second_request, serve) + build_client(features: { counting: feature }) + get + get + + assert_equal 2, @server.accepts + assert_equal %w[/ /], feature.requests + assert_equal 2, feature.responses + assert_equal 0, feature.errors + end + + def test_features_see_one_error_when_resend_fails + feature = CountingFeature.new + start(close_after_second_request, read_then_close) + build_client(features: { counting: feature }) + get + + assert_raises(HTTP::ConnectionError) { get } + assert_equal %w[/ /], feature.requests + assert_equal 1, feature.errors + end + + def test_resends_get_after_tls_connection_closes_without_close_notify + start(idle_close(:fin), serve, ssl: true) + build_ssl_client + get + @server.wait_closed + + assert_equal "ok", get + assert_equal 2, @server.accepts + end + + def test_resends_get_when_tls_connection_closes_after_reading_it + start(close_after_second_request, serve, ssl: true) + build_ssl_client + get + + assert_equal "ok", get + assert_equal 2, @server.accepts + end + + def test_resends_at_most_once_over_tls + start(close_after_second_request, read_then_close, ssl: true) + build_ssl_client + get + + # OpenSSL 3 reports an EOF without close_notify as SSLError, JRuby as a plain EOF + assert_raises(OpenSSL::SSL::SSLError, HTTP::ResponseHeaderError) { get } + assert_equal 3, @server.requests.size + assert_equal 2, @server.accepts + end + + def test_resends_get_after_tls_close_notify + start(idle_close(:close_notify), serve, ssl: true) + build_ssl_client + get + @server.wait_closed + + assert_equal "ok", get + assert_equal 2, @server.accepts + end + + def test_resends_get_through_proxy_after_tunnel_is_closed + start(serve, ssl: true) + @proxy = TunnelClosingProxy.new + Thread.new { @proxy.start } + @proxy.wait_ready + build_ssl_client(proxy: { proxy_address: @proxy.addr, proxy_port: @proxy.port }) + get + @proxy.close_tunnels + + assert_equal "ok", get + assert_equal 2, @proxy.connects + assert_equal 2, @server.accepts + end + + private + + def start(*handlers, ssl: false) + @server = ScriptedServer.new(*handlers, ssl: ssl) + end + + def build_client(**) + @client = HTTP::Client.new(persistent: @server.endpoint, **) + end + + def build_ssl_client(**) + build_client(ssl_context: SSLHelper.client_context, **) + end + + def get(path = "/") + @client.get("#{@server.endpoint}#{path}").to_s + end + + def serve + :serve.to_proc + end + + def answer(data) + lambda do |peer| + peer.read_request + peer.write(data) + end + end + + def read_then_close + lambda do |peer| + peer.read_request + peer.fin + end + end + + # Answers the first request, then closes the idle connection + def idle_close(action) + lambda do |peer| + peer.read_request + peer.respond + peer.public_send(action) + end + end + + # Answers the first request and hands the peer over once the second arrives + def on_second_request + lambda do |peer| + peer.read_request + peer.respond + peer.read_request + yield peer + end + end + + def close_after_second_request + on_second_request(&:fin) + end + + def reply_to_second_request(data) + on_second_request do |peer| + peer.write(data) + peer.fin + end + end +end diff --git a/test/http/connection_test.rb b/test/http/connection_test.rb index 8aae7f61..420a4aae 100644 --- a/test/http/connection_test.rb +++ b/test/http/connection_test.rb @@ -650,6 +650,64 @@ def test_send_request_with_small_body_pending_flushes_response_body assert flushed end + # --------------------------------------------------------------------------- + # #response_started? + # --------------------------------------------------------------------------- + def test_response_started_is_false_initially + refute_predicate build_connection, :response_started? + end + + def test_response_started_after_reading_headers + socket = fake(connect: nil, close: nil, readpartial: "HTTP/1.1 200 OK\r\nContent-Length: 0\r\n\r\n") + connection = build_connection(socket: socket) + connection.instance_variable_set(:@pending_response, true) + connection.read_headers! + + assert_predicate connection, :response_started? + end + + def test_response_not_started_when_connection_closes_before_any_byte + connection = build_connection(socket: fake(connect: nil, close: nil, readpartial: :eof)) + connection.instance_variable_set(:@pending_response, true) + + assert_raises(HTTP::ResponseHeaderError) { connection.read_headers! } + refute_predicate connection, :response_started? + end + + def test_response_started_when_first_bytes_fail_to_parse + connection = build_connection(socket: fake(connect: nil, close: nil, readpartial: "garbage\r\n\r\n")) + connection.instance_variable_set(:@pending_response, true) + + assert_raises(HTTP::ConnectionError) { connection.read_headers! } + assert_predicate connection, :response_started? + end + + def test_response_started_after_informational_response + responses = ["HTTP/1.1 100 Continue\r\n\r\n", :eof] + connection = build_connection(socket: fake(connect: nil, close: nil, readpartial: proc { responses.shift })) + connection.instance_variable_set(:@pending_response, true) + + assert_raises(HTTP::ResponseHeaderError) { connection.read_headers! } + assert_predicate connection, :response_started? + end + + def test_send_request_resets_response_started + connection = build_connection(socket: fake(connect: nil, close: nil, closed?: false, write: lambda(&:bytesize))) + connection.instance_variable_set(:@response_started, true) + connection.send_request(build_req) + + refute_predicate connection, :response_started? + end + + def test_send_request_resets_response_started_after_flushing_previous_response + connection = build_connection(socket: fake(connect: nil, close: nil, closed?: false, write: lambda(&:bytesize))) + response = fake(content_length: 2, flush: -> { connection.instance_variable_set(:@response_started, true) }) + connection.instance_variable_set(:@pending_response, response) + connection.send_request(build_req) + + refute_predicate connection, :response_started? + end + # --------------------------------------------------------------------------- # #finish_response # --------------------------------------------------------------------------- diff --git a/test/http/request_test.rb b/test/http/request_test.rb index 02720718..5fe920c7 100644 --- a/test/http/request_test.rb +++ b/test/http/request_test.rb @@ -353,6 +353,50 @@ def test_using_authenticated_proxy_with_four_keys_returns_true assert_predicate build_request(proxy: proxy), :using_authenticated_proxy? end + # #replayable? + + def test_replayable_with_idempotent_methods_returns_true + %i[get head options trace put delete].each do |verb| + assert_predicate build_request(verb: verb), :replayable?, verb + end + end + + def test_replayable_with_non_idempotent_methods_returns_false + %i[post patch connect].each do |verb| + refute_predicate build_request(verb: verb), :replayable?, verb + end + end + + def test_replayable_with_idempotency_key_returns_true + assert_predicate build_request(verb: :post, headers: { "Idempotency-Key" => "abc" }), :replayable? + end + + def test_replayable_with_x_idempotency_key_returns_true + assert_predicate build_request(verb: :post, headers: { "X-Idempotency-Key" => "abc" }), :replayable? + end + + def test_replayable_with_string_body_returns_true + assert_predicate build_request(verb: :put, body: "data"), :replayable? + end + + def test_replayable_with_string_subclass_body_returns_true + assert_predicate build_request(verb: :put, body: Class.new(String).new("data")), :replayable? + end + + def test_replayable_with_io_body_returns_false + refute_predicate build_request(verb: :put, body: StringIO.new("data")), :replayable? + end + + def test_replayable_with_enumerable_body_returns_false + refute_predicate build_request(verb: :put, body: ["data"]), :replayable? + end + + def test_replayable_with_idempotency_key_and_io_body_returns_false + request = build_request(verb: :post, headers: { "Idempotency-Key" => "abc" }, body: StringIO.new("data")) + + refute_predicate request, :replayable? + end + # #redirect def test_redirect_has_correct_uri diff --git a/test/support/http_handling_shared/connection_reuse_tests.rb b/test/support/http_handling_shared/connection_reuse_tests.rb index 852a39e1..420bc6d8 100644 --- a/test/support/http_handling_shared/connection_reuse_tests.rb +++ b/test/support/http_handling_shared/connection_reuse_tests.rb @@ -58,17 +58,25 @@ def test_connection_reuse_enabled_socket_issue_transparently_reopens first_socket_id = client.get("#{server.endpoint}/socket").body.to_s refute_equal "", first_socket_id - # Kill off the sockets we used - DummyServer::Servlet.sockets.each do |socket| - socket.close - rescue IOError - nil - end - DummyServer::Servlet.sockets.clear + kill_server_sockets - # Should error because we tried to use a bad socket + # A GET is resent once on a new socket + second_socket_id = client.get("#{server.endpoint}/socket").body.to_s + + refute_equal "", second_socket_id + refute_equal first_socket_id, second_socket_id + end + + def test_connection_reuse_enabled_socket_issue_raises_for_non_idempotent_request + client = build_client(persistent: server.endpoint) + first_socket_id = client.get("#{server.endpoint}/socket").body.to_s + + refute_equal "", first_socket_id + kill_server_sockets + + # Should error because a POST is not resent assert_raises(HTTP::ConnectionError) do - client.get("#{server.endpoint}/socket").body.to_s + client.post("#{server.endpoint}/sleep").body.to_s end # Should succeed since we create a new socket @@ -94,4 +102,15 @@ def test_connection_reuse_disabled_opens_new_sockets refute_includes sockets_used, "" assert_equal 2, sockets_used.uniq.length end + + private + + def kill_server_sockets + DummyServer::Servlet.sockets.each do |socket| + socket.close + rescue IOError + nil + end + DummyServer::Servlet.sockets.clear + end end diff --git a/test/support/scripted_server.rb b/test/support/scripted_server.rb new file mode 100644 index 00000000..d9e89360 --- /dev/null +++ b/test/support/scripted_server.rb @@ -0,0 +1,156 @@ +# frozen_string_literal: true + +require "socket" +require "openssl" + +require "support/proxy_server" +require "support/ssl_helper" + +# Serves each accepted connection with the next handler of a script +# +# Handlers receive a {Peer} and decide how the connection behaves: answer, +# close with or without a TLS close_notify, reset, or stall. Connections past +# the end of the script reuse the last handler. +class ScriptedServer + OK = "HTTP/1.1 200 OK\r\nContent-Length: 2\r\n\r\nok" + + # The server side of one accepted connection + class Peer + def initialize(server, socket) + @server = server + @socket = socket + end + + def read_request + head = +"" + while (line = @socket.gets) + head << line + break if line == "\r\n" + end + return if head.empty? + + length = head[/^Content-Length: (\d+)/i, 1].to_i + head << @socket.read(length) if length.positive? + @server.record(head) + end + + def write(data) + @socket.write(data) + end + + def respond + write(OK) + end + + def serve + respond while read_request + end + + # Closes the TCP socket without a TLS close_notify + def fin + @socket.to_io.close + @server.closed + end + + def close_notify + @socket.close + @server.closed + end + + def reset + @socket.to_io.setsockopt(Socket::SOL_SOCKET, Socket::SO_LINGER, [1, 0].pack("ii")) + fin + end + end + + attr_reader :accepts + + def initialize(*handlers, ssl: false) + @handlers = handlers + @ssl_context = SSLHelper.server_context if ssl + @tcp_server = TCPServer.new("127.0.0.1", 0) + @accepts = 0 + @requests = [] + @sockets = [] + @lock = Mutex.new + @closed = Queue.new + @thread = Thread.new { accept_loop } + end + + def endpoint + "#{@ssl_context ? 'https' : 'http'}://127.0.0.1:#{@tcp_server.addr[1]}" + end + + def requests + @lock.synchronize { @requests.dup } + end + + def record(request) + @lock.synchronize { @requests << request } + request + end + + def closed + @closed << true + end + + def wait_closed + @closed.pop(timeout: 5) || raise("no connection was closed") + end + + def shutdown + @tcp_server.close + @thread.join + @lock.synchronize { @sockets.each { |socket| socket.to_io.close unless socket.to_io.closed? } } + end + + private + + def accept_loop + loop do + socket = @tcp_server.accept + handler = @handlers[[@accepts, @handlers.size - 1].min] + @accepts += 1 + Thread.new { handle(socket, handler) } + end + rescue IOError, SystemCallError + nil + end + + def handle(socket, handler) + socket = upgrade(socket) if @ssl_context + @lock.synchronize { @sockets << socket } + handler.call(Peer.new(self, socket)) + rescue IOError, SystemCallError, OpenSSL::SSL::SSLError + nil + end + + def upgrade(socket) + ssl = OpenSSL::SSL::SSLSocket.new(socket, @ssl_context) + ssl.sync_close = true + ssl.accept + end +end + +# A proxy that can drop its open CONNECT tunnels the way an idle timeout does +class TunnelClosingProxy < ProxyServer + attr_reader :connects + + def initialize + super + @connects = 0 + @tunnels = Queue.new + end + + def close_tunnels + @tunnels.pop.close until @tunnels.empty? + end + + private + + def tunnel_connection(client, target) + @connects += 1 + @tunnels << client + super + end +end