From bda46a92d2fd72ba78a93cd376bfa746ffbae614 Mon Sep 17 00:00:00 2001 From: ilyazub Date: Thu, 1 Oct 2026 15:34:15 +0200 Subject: [PATCH] Reconnect instead of reusing idle connections the server closed A server or proxy can close a keep-alive connection while it sits idle in a persistent client. The client expires idle connections after keep_alive_timeout, 5 seconds by default, but a server's idle timeout can be shorter, and the server starts counting when it sends the response, before the client has read it. The client only checked whether its own end of the socket was closed, so it wrote the next request into the closed connection. That request failed with HTTP::ResponseHeaderError ("couldn't read response headers"), or with OpenSSL::SSL::SSLError ("unexpected eof while reading") when the server skipped close_notify. A server that sends a response before closing, such as a 408 Request Timeout, was worse: the next request read that 408 as its own response. Before reusing a connection, read off the previous response, then check the socket without blocking. Once the previous response has been read, an idle connection has nothing left to read, so any readable data means it can't carry another request: an EOF, a reset, a TLS close_notify or an unsolicited response. Close it then and open a new one. urllib3 makes the same check (wait_for_read(sock, timeout=0.0)), Net::HTTP reconnects on wait_readable(0) && eof?, and Go's transport drops idle connections that see an EOF or an unsolicited response. www.debian.org, lwn.net, www.postgresql.org and www.sqlite.org closed idle connections 4.8 to 4.98 seconds after the client read a response, and www.haproxy.org as early as 2 seconds. With the default keep_alive_timeout, a second request 4.99 seconds later (3 seconds for www.haproxy.org) failed 16 of 16 times and now succeeds 16 of 16 times on a new connection. 14 servers that keep idle connections open were still reused 28 of 28 times after 4.5 seconds idle. The check happens before any byte of the next request is written, so it applies to every method, POST included, and can't send a request twice. It runs after the client is marked dirty, so an exception that interrupts reading off the previous response still makes the next request reconnect. It can't help when the server closes the connection after the request was written; only idempotent requests are safe to retry then. Reading off a body larger than 1 MiB closes the connection instead, but the client then wrote the next request to the closed connection and raised HTTP::SocketWriteError ("closed stream"). The same check now reconnects. test_connection_reuse_enabled_socket_issue_transparently_reopens asserted that the first request after the server closed the socket raised. It now expects the client to reopen the connection, as its name says, and a POST variant covers the same case. --- CHANGELOG.md | 19 ++++ lib/http/client.rb | 18 ++++ lib/http/connection.rb | 2 +- lib/http/connection/internals.rb | 39 +++++++- sig/http.rbs | 5 +- test/http/client_test.rb | 80 ++++++++++++++- test/http/connection_test.rb | 99 +++++++++++++++++++ test/support/dummy_server/routes.rb | 9 ++ .../connection_reuse_tests.rb | 76 +++++++++++--- 9 files changed, 328 insertions(+), 19 deletions(-) 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