diff --git a/CHANGELOG.md b/CHANGELOG.md index 954aa1ae..9ecfeead 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,8 +7,25 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +### Added + +- `HTTP::Connection#stale?` checks, without blocking, whether the server + closed an idle connection or sent data on it. + `HTTP::Connection#flush_pending_response` is now public. + ### Fixed +- Persistent connections are no longer reused after the server or a proxy + closed them while idle. The next request used to fail with + `HTTP::ResponseHeaderError` ("couldn't read response headers") or + `OpenSSL::SSL::SSLError` ("unexpected eof while reading"), or read an + unsolicited response the server sent before closing, such as + `408 Request Timeout`, as its own. The client now checks the idle socket + before reuse and reconnects when anything is readable. ([#420], [#459]) +- Reusing a persistent connection after leaving a response body larger than + 1 MiB unread no longer raises `HTTP::SocketWriteError` ("closed stream"). + The connection is closed to skip the body, and the client now reconnects + instead of writing to it. - 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 +315,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..3538881f 100644 --- a/lib/http/client.rb +++ b/lib/http/client.rb @@ -144,6 +144,7 @@ def perform_with_retry(req, options) def send_request(req, options) notify_features(req, options) + discard_stale_connection @connection ||= HTTP::Connection.new(req, options) unless @connection.failed_proxy_connect? @@ -155,6 +156,23 @@ def send_request(req, options) raise end + # Drop the connection unless it can carry another request + # + # Reads off the previous response first. Runs after the client is marked + # dirty, so an interrupted read still makes the next request reconnect. + # + # @return [void] + # @api private + def discard_stale_connection + return unless @connection + + @connection.flush_pending_response + return if @connection.keep_alive? && !@connection.stale? + + @connection.close + @connection = nil + end + # Build response and apply feature wrapping # @return [HTTP::Response] the wrapped response # @api private diff --git a/lib/http/connection.rb b/lib/http/connection.rb index e550a4b6..d8672a5e 100644 --- a/lib/http/connection.rb +++ b/lib/http/connection.rb @@ -101,7 +101,7 @@ def failed_proxy_connect? # @return [nil] # @api public def send_request(req) - flush_pending_response if @pending_response + flush_pending_response if @pending_request raise StateError, "Tried to send a request while a response is pending. Make sure you read off the body." diff --git a/lib/http/connection/internals.rb b/lib/http/connection/internals.rb index 3cb662eb..3f246797 100644 --- a/lib/http/connection/internals.rb +++ b/lib/http/connection/internals.rb @@ -2,14 +2,45 @@ module HTTP class Connection - # Internal private methods for Connection + # Lower-level socket and response handling for Connection module Internals - private + # Whether the server closed this idle connection or sent data on it + # + # Checks the socket without blocking. On an idle connection any readable + # data means it can't carry another request: an EOF, a reset, a TLS + # close_notify, or a response nobody asked for, such as the 408 some + # servers send before closing. Always false while a response is pending, + # because its unread body is expected data. + # + # @example + # connection.stale? + # + # @return [Boolean] + # @api public + def stale? + return false if @pending_response + + io = @socket.socket if @socket.respond_to?(:socket) + return false unless io.respond_to?(:to_io) + + io.to_io.wait_readable(0) ? true : false + rescue IOError, SystemCallError + true + end # Flush the pending response body so the connection can be reused + # + # Closes the connection instead when the body can't be flushed or is + # larger than {MAX_FLUSH_SIZE}. Does nothing when no response is pending. + # + # @example + # connection.flush_pending_response + # # @return [void] - # @api private + # @api public def flush_pending_response + return unless @pending_response + response = @pending_response unless response.respond_to?(:flush) close @@ -21,6 +52,8 @@ def flush_pending_response close end + private + # Flush the response or close if the body exceeds the size limit # @param response [HTTP::Response] the response to flush # @return [void] diff --git a/sig/http.rbs b/sig/http.rbs index 9b1c4a69..006beb90 100644 --- a/sig/http.rbs +++ b/sig/http.rbs @@ -172,6 +172,7 @@ module HTTP 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 discard_stale_connection: () -> void def verify_connection!: (URI uri) -> void end @@ -1170,9 +1171,11 @@ module HTTP def close: () -> void end + def stale?: () -> bool + def flush_pending_response: () -> void + private - def flush_pending_response: () -> void def flush_or_close_response: (Response response) -> void def start_tls: (Request req, Options options) -> void def send_proxy_connect_request: (Request req) -> void diff --git a/test/http/client_test.rb b/test/http/client_test.rb index 8e8d4d36..dcf6e639 100644 --- a/test/http/client_test.rb +++ b/test/http/client_test.rb @@ -682,13 +682,87 @@ def test_perform_with_failed_proxy_connect_skips_sending_request close: nil, "pending_response=": ->(*) {} ) - proxy_client.instance_variable_set(:@connection, conn) - proxy_client.instance_variable_set(:@state, :clean) req = HTTP::Request.new(verb: :get, uri: "http://example.com/", headers: {}) - response = proxy_client.perform(req, HTTP::Options.new) + response = HTTP::Connection.stub(:new, conn) { proxy_client.perform(req, HTTP::Options.new) } assert_equal 407, response.status.to_i end + + # #perform on a persistent connection + + def build_idle_connection(**overrides) + fake( + failed_proxy_connect?: false, + send_request: nil, + read_headers!: nil, + proxy_response_headers: {}, + status_code: 200, + http_version: "1.1", + headers: HTTP::Headers.new, + finish_response: nil, + keep_alive?: true, + expired?: false, + flush_pending_response: nil, + stale?: false, + close: nil, + "pending_response=": ->(*) {}, + **overrides + ) + end + + def perform_over(connection) + client = HTTP::Client.new + client.instance_variable_set(:@connection, connection) + client.instance_variable_set(:@state, :clean) + replacement = build_idle_connection + req = HTTP::Request.new(verb: :get, uri: "http://example.com/", headers: {}) + HTTP::Connection.stub(:new, replacement) { client.perform(req, HTTP::Options.new) } + + [client.instance_variable_get(:@connection), replacement] + end + + def test_perform_reuses_idle_connection + connection = build_idle_connection + + assert_same connection, perform_over(connection).first + end + + def test_perform_replaces_stale_connection + closed = false + connection = build_idle_connection(stale?: true, close: -> { closed = true }) + current, replacement = perform_over(connection) + + assert closed + assert_same replacement, current + end + + def test_perform_replaces_connection_closed_while_flushing_previous_response + alive = true + closed = false + connection = build_idle_connection( + keep_alive?: -> { alive }, + flush_pending_response: -> { alive = false }, + close: -> { closed = true } + ) + current, replacement = perform_over(connection) + + assert closed + assert_same replacement, current + end + + def test_perform_keeps_client_dirty_when_interrupted_on_replacement_connection + client = HTTP::Client.new + client.instance_variable_set(:@connection, build_idle_connection(stale?: true)) + client.instance_variable_set(:@state, :clean) + interrupted = build_idle_connection(send_request: ->(*) { raise Interrupt }) + req = HTTP::Request.new(verb: :get, uri: "http://example.com/", headers: {}) + + HTTP::Connection.stub(:new, interrupted) do + assert_raises(Interrupt) { client.perform(req, HTTP::Options.new) } + end + + assert_equal :dirty, client.instance_variable_get(:@state) + end end class HTTPClientHTTPHandlingTest < Minitest::Test diff --git a/test/http/connection_test.rb b/test/http/connection_test.rb index 8aae7f61..3102f7da 100644 --- a/test/http/connection_test.rb +++ b/test/http/connection_test.rb @@ -978,6 +978,105 @@ def test_expired_returns_false_when_connection_has_not_expired refute_predicate connection, :expired? end + # --------------------------------------------------------------------------- + # #stale? + # --------------------------------------------------------------------------- + def build_connection_over(io) + build_connection(socket: fake(connect: nil, close: nil, closed?: false, socket: io)) + end + + def test_stale_returns_true_when_idle_socket_is_readable + raw = Object.new + connection = build_connection_over(fake(to_io: fake(wait_readable: raw))) + + assert_same true, connection.stale? + end + + def test_stale_returns_false_when_idle_socket_is_not_readable + connection = build_connection_over(fake(to_io: fake(wait_readable: nil))) + + assert_same false, connection.stale? + end + + def test_stale_checks_readability_without_waiting + timeouts = [] + raw = fake(wait_readable: ->(timeout) { timeouts << timeout and nil }) + build_connection_over(fake(to_io: raw)).stale? + + assert_equal [0], timeouts + end + + def test_stale_returns_false_while_response_is_pending + connection = build_connection_over(fake(to_io: fake(wait_readable: Object.new))) + connection.instance_variable_set(:@pending_response, true) + + assert_same false, connection.stale? + end + + def test_stale_returns_false_when_socket_is_not_exposed + connection = build_connection(socket: fake(connect: nil, close: nil, closed?: false)) + + assert_same false, connection.stale? + end + + def test_stale_returns_false_when_socket_is_not_an_io + connection = build_connection_over(Object.new) + + assert_same false, connection.stale? + end + + def test_stale_returns_true_when_socket_is_closed + raw = fake(wait_readable: ->(_) { raise IOError, "closed stream" }) + connection = build_connection_over(fake(to_io: raw)) + + assert_same true, connection.stale? + end + + def test_stale_returns_true_when_socket_errors + raw = fake(wait_readable: ->(_) { raise Errno::EBADF }) + connection = build_connection_over(fake(to_io: raw)) + + assert_same true, connection.stale? + end + + # --------------------------------------------------------------------------- + # #flush_pending_response + # --------------------------------------------------------------------------- + def test_flush_pending_response_does_nothing_without_pending_response + closed = false + connection = build_connection(socket: fake(connect: nil, close: -> { closed = true }, closed?: false)) + connection.flush_pending_response + + refute closed + end + + def test_flush_pending_response_flushes_pending_response + connection = build_connection + flushed = false + connection.instance_variable_set(:@pending_response, fake(content_length: 1, flush: -> { flushed = true })) + connection.flush_pending_response + + assert flushed + end + + def test_flush_pending_response_closes_when_response_cannot_be_flushed + closed = false + connection = build_connection(socket: fake(connect: nil, close: -> { closed = true }, closed?: false)) + connection.instance_variable_set(:@pending_response, true) + connection.flush_pending_response + + assert closed + end + + def test_flush_pending_response_closes_when_flushing_fails + closed = false + connection = build_connection(socket: fake(connect: nil, close: -> { closed = true }, closed?: false)) + connection.instance_variable_set(:@pending_response, fake(content_length: nil, flush: -> { raise IOError })) + connection.flush_pending_response + + assert closed + end + # --------------------------------------------------------------------------- # keep_alive behavior (set_keep_alive) # --------------------------------------------------------------------------- diff --git a/test/support/dummy_server/routes.rb b/test/support/dummy_server/routes.rb index ed13ec26..a6246cc6 100644 --- a/test/support/dummy_server/routes.rb +++ b/test/support/dummy_server/routes.rb @@ -111,6 +111,15 @@ class Servlet res.body = bytes.pack("c*") end + get "/large" do |_req, res| + res.status = 200 + res.body = "x" * (HTTP::Connection::MAX_FLUSH_SIZE + 1) + end + + get "/close" do |req, _res| + req.socket.close + end + get "/iso-8859-1" do |_req, res| res["Content-Type"] = "text/plain; charset=ISO-8859-1" res.body = "testæ".encode(Encoding::ISO8859_1) diff --git a/test/support/http_handling_shared/connection_reuse_tests.rb b/test/support/http_handling_shared/connection_reuse_tests.rb index 852a39e1..18854569 100644 --- a/test/support/http_handling_shared/connection_reuse_tests.rb +++ b/test/support/http_handling_shared/connection_reuse_tests.rb @@ -56,25 +56,59 @@ def test_connection_reuse_enabled_reading_cached_body_succeeds def test_connection_reuse_enabled_socket_issue_transparently_reopens client = build_client(persistent: server.endpoint) first_socket_id = client.get("#{server.endpoint}/socket").body.to_s + client_socket = idle_client_socket(client) refute_equal "", first_socket_id - # Kill off the sockets we used + kill_server_sockets + wait_for_server_bytes(client_socket) + + second_socket_id = client.get("#{server.endpoint}/socket").body.to_s + + refute_equal first_socket_id, second_socket_id + assert_predicate client_socket, :closed? + end + + def test_connection_reuse_enabled_socket_issue_transparently_reopens_for_post + client = build_client(persistent: server.endpoint) + client.get("#{server.endpoint}/socket").body.to_s + client_socket = idle_client_socket(client) + + kill_server_sockets + wait_for_server_bytes(client_socket) + + assert_equal "hello", client.post("#{server.endpoint}/sleep").body.to_s + assert_predicate client_socket, :closed? + end + + def test_connection_reuse_enabled_reopens_when_server_responds_while_idle + client = build_client(persistent: server.endpoint) + first_socket_id = client.get("#{server.endpoint}/socket").body.to_s + client_socket = idle_client_socket(client) + DummyServer::Servlet.sockets.each do |socket| - socket.close - rescue IOError - nil + socket.write("HTTP/1.1 408 Request Timeout\r\nContent-Length: 0\r\n\r\n") end DummyServer::Servlet.sockets.clear + wait_for_server_bytes(client_socket) - # Should error because we tried to use a bad socket - assert_raises(HTTP::ConnectionError) do - client.get("#{server.endpoint}/socket").body.to_s - end + response = client.get("#{server.endpoint}/socket") - # Should succeed since we create a new socket - second_socket_id = client.get("#{server.endpoint}/socket").body.to_s + assert_equal 200, response.code + refute_equal first_socket_id, response.body.to_s + end - refute_equal first_socket_id, second_socket_id + def test_connection_reuse_enabled_reopens_when_unread_body_is_too_large_to_flush + client = build_client(persistent: server.endpoint) + client.get("#{server.endpoint}/large") + + assert_equal "", client.get(server.endpoint).body.to_s + end + + def test_connection_reuse_enabled_raises_when_server_closes_after_receiving_request + client = build_client(persistent: server.endpoint) + client.get("#{server.endpoint}/socket").body.to_s + + assert_raises(HTTP::ConnectionError) { client.get("#{server.endpoint}/close") } end def test_connection_reuse_enabled_change_in_host_errors @@ -94,4 +128,24 @@ def test_connection_reuse_disabled_opens_new_sockets refute_includes sockets_used, "" assert_equal 2, sockets_used.uniq.length end + + private + + def idle_client_socket(client) + client.instance_variable_get(:@connection).instance_variable_get(:@socket).socket.to_io + end + + # Loopback delivers a close or write asynchronously; wait until it lands + def wait_for_server_bytes(socket) + assert socket.wait_readable(5), "server bytes never reached the client socket" + end + + def kill_server_sockets + DummyServer::Servlet.sockets.each do |socket| + socket.close + rescue IOError + nil + end + DummyServer::Servlet.sockets.clear + end end