diff --git a/CHANGELOG.md b/CHANGELOG.md index 9f3cc55a..8009c29c 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,14 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +### Fixed + +- (backported) Reconnect instead of reusing a persistent connection the server + closed or wrote to while it sat idle, which failed the next request with + "couldn't read response headers" or returned the server's stray 408 + ([#420](https://github.com/httprb/http/issues/420), + [#459](https://github.com/httprb/http/issues/459)) + ## [5.3.1] - 2025-06-09 diff --git a/lib/http/client.rb b/lib/http/client.rb index 3ecb184a..9502a7bc 100644 --- a/lib/http/client.rb +++ b/lib/http/client.rb @@ -128,7 +128,8 @@ def verify_connection!(uri) # We re-create the connection object because we want to let prior requests # lazily load the body as long as possible, and this mimics prior functionality. - return close if @connection && (!@connection.keep_alive? || @connection.expired?) + # A connection the server closed while idle still looks open until we read from it. + return close if @connection && (!@connection.keep_alive? || @connection.expired? || @connection.stale?) # If we get into a bad state (eg, Timeout.timeout ensure being killed) # close the connection to prevent potential for mixed responses. diff --git a/lib/http/connection.rb b/lib/http/connection.rb index adeac882..19b50547 100644 --- a/lib/http/connection.rb +++ b/lib/http/connection.rb @@ -152,6 +152,26 @@ def expired? !@conn_expires_at || @conn_expires_at < Time.now end + # 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. + # + # @return [Boolean] + 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 + private # Sets up SSL context and starts TLS if needed. diff --git a/spec/lib/http/client_spec.rb b/spec/lib/http/client_spec.rb index 88822f3d..5e0fe9f2 100644 --- a/spec/lib/http/client_spec.rb +++ b/spec/lib/http/client_spec.rb @@ -441,6 +441,81 @@ def on_error(request, error) end end + # Self-contained because "working with SSL" is disabled (#627) + describe "reusing a persistent TLS connection the server closed" do + let(:key) { OpenSSL::PKey::RSA.new(2048) } + let(:cert) do + OpenSSL::X509::Certificate.new.tap do |cert| + cert.version = 2 + cert.serial = 1 + cert.subject = OpenSSL::X509::Name.parse("/CN=127.0.0.1") + cert.issuer = cert.subject + cert.public_key = key.public_key + cert.not_before = Time.now - 60 + cert.not_after = Time.now + 3600 + extensions = OpenSSL::X509::ExtensionFactory.new(cert, cert) + cert.add_extension(extensions.create_extension("subjectAltName", "IP:127.0.0.1")) + cert.sign(key, OpenSSL::Digest.new("SHA256")) + end + end + let(:tcp_server) { TCPServer.new("127.0.0.1", 0) } + let(:endpoint) { "https://127.0.0.1:#{tcp_server.addr[1]}" } + let(:server_connections) { Queue.new } + + let(:client) do + context = OpenSSL::SSL::SSLContext.new + context.verify_mode = OpenSSL::SSL::VERIFY_PEER + context.cert_store = OpenSSL::X509::Store.new.tap { |store| store.add_cert(cert) } + described_class.new(:persistent => endpoint, :ssl_context => context) + end + + # Answers each request with the number of the connection it arrived on + let!(:server_thread) do + context = OpenSSL::SSL::SSLContext.new + context.cert = cert + context.key = key + tls_server = OpenSSL::SSL::SSLServer.new(tcp_server, context) + + Thread.new do + (1..).each do |number| + connection = tls_server.accept + server_connections << connection + Thread.new { serve(connection, number.to_s) } + end + rescue IOError, SystemCallError, OpenSSL::SSL::SSLError + nil + end + end + + after do + tcp_server.close + server_thread.join(1) + end + + def serve(connection, body) + while connection.gets + loop do + header = connection.gets + break if header.nil? || header == "\r\n" + end + connection.write("HTTP/1.1 200 OK\r\nContent-Length: #{body.bytesize}\r\n\r\n#{body}") + end + rescue IOError, SystemCallError, OpenSSL::SSL::SSLError + nil + end + + it "reconnects instead of sending the request on it" do + expect(client.get("#{endpoint}/").body.to_s).to eq("1") + client_socket = client.instance_variable_get(:@connection).instance_variable_get(:@socket).socket.to_io + + server_connections.pop.close + expect(client_socket.wait_readable(5)).to be_truthy + + expect(client.get("#{endpoint}/").body.to_s).to eq("2") + expect(client_socket).to be_closed + end + end + describe "#perform" do let(:client) { described_class.new } diff --git a/spec/lib/http/connection_spec.rb b/spec/lib/http/connection_spec.rb index 719cfdca..78e92d66 100644 --- a/spec/lib/http/connection_spec.rb +++ b/spec/lib/http/connection_spec.rb @@ -85,4 +85,54 @@ expect(connection.finished_request?).to be true end end + + describe "#stale?" do + let(:pair) { UNIXSocket.pair } + let(:io) { pair.first } + let(:peer) { pair.last } + let(:socket) { double(:connect => nil, :close => nil, :socket => io) } + + after { pair.each { |s| s.close unless s.closed? } } + + it "is false while the idle peer is still there" do + expect(connection.stale?).to be false + end + + it "is true once the peer closed the idle connection" do + peer.close + expect(connection.stale?).to be true + end + + it "is true once the peer sent data on the idle connection" do + peer.write("HTTP/1.1 408 Request Timeout\r\nContent-Length: 0\r\n\r\n") + expect(connection.stale?).to be true + end + + it "is false while a response is pending, because its body is expected data" do + peer.write("unread body") + connection.instance_variable_set(:@pending_response, true) + expect(connection.stale?).to be false + end + + it "is true when the socket is already closed" do + io.close + expect(connection.stale?).to be true + end + + context "when the timeout class doesn't expose its socket" do + let(:socket) { double(:connect => nil, :close => nil) } + + it "is false" do + expect(connection.stale?).to be false + end + end + + context "when the exposed socket isn't an IO" do + let(:socket) { double(:connect => nil, :close => nil, :socket => Object.new) } + + it "is false" do + expect(connection.stale?).to be false + end + end + end end diff --git a/spec/support/dummy_server/servlet.rb b/spec/support/dummy_server/servlet.rb index 075f8ad3..8cd2d7c0 100644 --- a/spec/support/dummy_server/servlet.rb +++ b/spec/support/dummy_server/servlet.rb @@ -72,6 +72,10 @@ def do_#{method.upcase}(req, res) end end + get "/close" do |req, _res| + req.instance_variable_get(:@socket).close + end + get "/params" do |req, res| next not_found(req, res) unless "foo=bar" == req.query_string diff --git a/spec/support/http_handling_shared.rb b/spec/support/http_handling_shared.rb index ff13b0a3..b331182c 100644 --- a/spec/support/http_handling_shared.rb +++ b/spec/support/http_handling_shared.rb @@ -154,20 +154,71 @@ it "transparently reopens", :flaky do first_socket_id = client.get("#{server.endpoint}/socket").body.to_s expect(first_socket_id).to_not eq("") - # Kill off the sockets we used - # rubocop:disable Style/RescueModifier + client_socket = idle_client_socket(client) + + kill_server_sockets + wait_for_server_bytes(client_socket) + + second_socket_id = client.get("#{server.endpoint}/socket").body.to_s + expect(second_socket_id).to_not eq(first_socket_id) + expect(client_socket).to be_closed + end + + it "transparently reopens for a POST", :flaky do + client.get("#{server.endpoint}/socket").body.to_s + client_socket = idle_client_socket(client) + + kill_server_sockets + wait_for_server_bytes(client_socket) + + expect(client.post("#{server.endpoint}/echo-body", :body => "sent once").body.to_s).to eq("sent once") + expect(client_socket).to be_closed + end + + [ + [HTTP::Timeout::PerOperation, {:connect_timeout => 5, :read_timeout => 5, :write_timeout => 5}], + [HTTP::Timeout::Global, {:global_timeout => 5}] + ].each do |timeout_class, timeout_options| + context "with #{timeout_class}" do + let(:extra_options) { {:timeout_class => timeout_class, :timeout_options => timeout_options} } + + it "transparently reopens", :flaky do + first_socket_id = client.get("#{server.endpoint}/socket").body.to_s + client_socket = idle_client_socket(client) + + kill_server_sockets + wait_for_server_bytes(client_socket) + + expect(client.get("#{server.endpoint}/socket").body.to_s).to_not eq(first_socket_id) + expect(client_socket).to be_closed + end + end + end + end + + context "when the server responds while the connection is idle" do + it "reopens instead of reading that response", :flaky do + DummyServer::Servlet.sockets.clear + 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 nil + socket.write("HTTP/1.1 408 Request Timeout\r\nContent-Length: 0\r\n\r\n") end DummyServer::Servlet.sockets.clear - # rubocop:enable Style/RescueModifier + wait_for_server_bytes(client_socket) - # Should error because we tried to use a bad socket - expect { client.get("#{server.endpoint}/socket").body.to_s }.to raise_error HTTP::ConnectionError + response = client.get("#{server.endpoint}/socket") + expect(response.code).to eq(200) + expect(response.body.to_s).to_not eq(first_socket_id) + end + end - # Should succeed since we create a new socket - second_socket_id = client.get("#{server.endpoint}/socket").body.to_s - expect(second_socket_id).to_not eq(first_socket_id) + context "when the server closes the connection after receiving the request" do + it "raises" do + client.get("#{server.endpoint}/socket").body.to_s + + expect { client.get("#{server.endpoint}/close") }.to raise_error(HTTP::ConnectionError) end end @@ -187,4 +238,22 @@ end end end + + 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) + expect(socket.wait_readable(5)).to be_truthy, "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