diff --git a/lib/async/http/body/pipe.rb b/lib/async/http/body/pipe.rb index da221778..158cdb27 100644 --- a/lib/async/http/body/pipe.rb +++ b/lib/async/http/body/pipe.rb @@ -33,6 +33,13 @@ def to_io @tail end + # Close the socket and stop reading from the input. Unlike {close}, the writer forwards any data already written to the socket before closing the output normally. + def finish + @reader&.stop + + @tail.close + end + # Close the pipe and stop the reader and writer tasks. def close @reader&.stop diff --git a/lib/async/http/proxy.rb b/lib/async/http/proxy.rb index 4f295f7e..f42151fe 100644 --- a/lib/async/http/proxy.rb +++ b/lib/async/http/proxy.rb @@ -15,14 +15,50 @@ module HTTP class Proxy # Raised when a CONNECT tunnel through a proxy cannot be established. class ConnectFailure < StandardError + # The maximum number of bytes of the response body to read, see {body}. + MAXIMUM_BODY_SIZE = 1024 * 8 + # Initialize the failure with the unsuccessful response. # @parameter response [Protocol::HTTP::Response] The failed response from the proxy. - def initialize(response) + # @parameter body [String | Nil] The beginning of the response body, if any. + def initialize(response, body = nil) super "Failed to connect: #{response.status}" @response = response + @body = body end + # @attribute [Protocol::HTTP::Response] The failed response from the proxy. It has already been closed. attr :response + + # @attribute [String | Nil] Up to {MAXIMUM_BODY_SIZE} bytes of the response body, which may explain why the proxy rejected the connection. + # + # The body is read before this error is raised, so a proxy which neither finishes the body nor closes the connection (e.g. an HTTP/1 response without `content-length` on a persistent connection) will delay the failure until the caller's own timeout, if any. + attr :body + end + + # Extends the socket returned by {Proxy#connect} so that closing it closes the entire tunnel. + # + # Closing the socket of a {Body::Pipe} only ends the CONNECT request body, so the CONNECT response would remain open (and the proxy connection busy) until the proxy closes its end of the tunnel. Instead, closing this socket finishes the pipe: the CONNECT response is closed immediately, and the CONNECT request is closed once any data already written to the socket has been forwarded. HTTP/2 then resets the stream and HTTP/1 closes the connection, neither of which can be reused after a tunnel anyway. + module Tunnel + # @parameter pipe [Body::Pipe] The pipe connected to the CONNECT request and response. + # @returns [Socket] The socket for the pipe, which finishes the pipe when closed. + def self.wrap(pipe) + socket = pipe.to_io + socket.extend(self) + socket.instance_variable_set(:@pipe, pipe) + + return socket + end + + # Finish the pipe, which also closes the socket. + def close + if pipe = @pipe + @pipe = nil + pipe.finish + else + super + end + end end # Extends {Async::HTTP::Client} with proxy capabilities. @@ -99,19 +135,24 @@ def connect(&block) if response.success? pipe = Body::Pipe.new(response.body, input) - return pipe.to_io unless block_given? + return Tunnel.wrap(pipe) unless block_given? begin yield pipe.to_io ensure - pipe.close + pipe.finish end else # This ensures we don't leave a response dangling: input.close - response.close - raise ConnectFailure, response + begin + body = read_failure_body(response) + ensure + response.close + end + + raise ConnectFailure.new(response, body) end end @@ -119,6 +160,27 @@ def connect(&block) def wrap_endpoint(endpoint) Endpoint.new(endpoint.url, self, **endpoint.options) end + + private + + # Proxies often explain why they rejected a connection in the response body, so read the beginning of it. Failing to do so shouldn't hide the original failure. + def read_failure_body(response) + body = response.body + return nil unless body + + buffer = String.new + + begin + while buffer.bytesize < ConnectFailure::MAXIMUM_BODY_SIZE and chunk = body.read + buffer << chunk + end + rescue IOError, SystemCallError, ::Protocol::HTTP::Error => error + # Anything else, e.g. `Async::TimeoutError`, should propagate to the caller: + Console.debug(self, "Failed to read proxy response body!", error) + end + + return buffer.byteslice(0, ConnectFailure::MAXIMUM_BODY_SIZE) + end end Client.prepend(Proxy::Client) diff --git a/releases.md b/releases.md index f56fb22f..76e01a2c 100644 --- a/releases.md +++ b/releases.md @@ -1,5 +1,10 @@ # Releases +## Unreleased + + - Closing the socket returned by `Async::HTTP::Proxy#connect` (without a block) now closes the entire CONNECT tunnel, rather than only ending the CONNECT request body. The CONNECT response is closed immediately, and the CONNECT request is closed once any data already written to the socket has been forwarded, after which HTTP/2 proxies receive `RST_STREAM(NO_ERROR)` and HTTP/1 proxy connections are closed. `Async::HTTP::Client#close` for the proxy client no longer waits for the proxy to close its end of the tunnel. The block form of `Async::HTTP::Proxy#connect` closes the tunnel the same way when the block exits, and no longer discards data written immediately before the block exits. + - `Async::HTTP::Proxy::ConnectFailure#body` exposes up to 8 KiB of the proxy's response body, which often explains why the connection was rejected. I/O and protocol errors while reading it are ignored, but other errors, such as `Async::TimeoutError`, propagate to the caller. + ## v0.105.0 - Respond with `400 Bad Request` for any `Protocol::HTTP::BadRequest` raised while parsing HTTP/1 or HTTP/2 requests, include the exception class name without reflecting request data, and avoid reporting them as unhandled server errors. diff --git a/test/async/http/body/pipe.rb b/test/async/http/body/pipe.rb index 3d043b0a..1542e394 100644 --- a/test/async/http/body/pipe.rb +++ b/test/async/http/body/pipe.rb @@ -54,6 +54,49 @@ def before end end + with "#finish" do + include Sus::Fixtures::Async::ReactorContext + + let(:output) {Async::HTTP::Body::Writable.new} + let(:pipe) {subject.new(input, output)} + + it "forwards data written to the socket before closing the output" do + pipe.to_io.write(data) + pipe.finish + + expect(pipe.to_io).to be(:closed?) + expect(input).to be(:closed?) + + expect(output.read).to be == data + expect(output.read).to be_nil + end + end + + with "#close" do + include Sus::Fixtures::Async::ReactorContext + + let(:output) {Async::HTTP::Body::Writable.new} + let(:pipe) {subject.new(input, output)} + + it "closes the input when forwarding to the closed socket fails" do + pipe.to_io.close + + # Closing the socket only ends the output: + expect(output.read).to be_nil + expect(input).not.to be(:closed?) + + input.write(data) + + Async::Task.current.with_timeout(1) do + Async::Task.current.yield until input.closed? + end + + expect(input).to be(:closed?) + ensure + pipe.close + end + end + with "reactor going out of scope" do it "finishes" do # ensures pipe background tasks are transient diff --git a/test/async/http/proxy.rb b/test/async/http/proxy.rb index 34b99df3..2f9c8242 100644 --- a/test/async/http/proxy.rb +++ b/test/async/http/proxy.rb @@ -139,14 +139,13 @@ end end - it "closes the response when forwarding to the closed peer fails" do + it "keeps the response open when the request is closed" do proxy = Async::HTTP::Proxy.tcp(client, "localhost", 1) peer = proxy.connect expect(proxy.client.pool).to be(:busy?) - peer.close - peer = nil + peer.close_write current_task = Async::Task.current current_task.with_timeout(1) do @@ -157,6 +156,8 @@ write_response.resolve(true) + expect(peer.read).to be == "Hello World!" + current_task.with_timeout(1) do current_task.yield while proxy.client.pool.busy? end @@ -167,6 +168,244 @@ peer&.close proxy&.close end + + it "closes the response when the peer is closed" do + proxy = Async::HTTP::Proxy.tcp(client, "localhost", 1) + peer = proxy.connect + + expect(proxy.client.pool).to be(:busy?) + + peer.close + expect(peer).to be(:closed?) + + current_task = Async::Task.current + current_task.with_timeout(1) do + request_closed.wait + current_task.yield while proxy.client.pool.busy? + end + + # The remote end never finished the response, but the tunnel was closed anyway: + expect(write_response).not.to be(:resolved?) + expect(proxy.client.pool).not.to be(:busy?) + ensure + peer&.close + proxy&.close + end + end + + with "rejected CONNECT" do + let(:reason) {"Egress proxying is denied to host 'localhost'."} + + let(:app) do + Protocol::HTTP::Middleware.for do |request| + Protocol::HTTP::Response[403, {"content-type" => "text/plain"}, [reason]] + end + end + + it "exposes the response body" do + proxy = Async::HTTP::Proxy.tcp(client, "localhost", 1) + + expect do + proxy.connect + end.to raise_exception(Async::HTTP::Proxy::ConnectFailure, message: be == "Failed to connect: 403") + + begin + proxy.connect + rescue Async::HTTP::Proxy::ConnectFailure => error + expect(error.response.status).to be == 403 + expect(error.body).to be == reason + end + + expect(error).to be_a(Async::HTTP::Proxy::ConnectFailure) + + current_task = Async::Task.current + current_task.with_timeout(1) do + current_task.yield while proxy.client.pool.busy? + end + ensure + proxy&.close + end + + with "large response body" do + let(:reason) {"x" * (Async::HTTP::Proxy::ConnectFailure::MAXIMUM_BODY_SIZE * 4)} + + it "reads a bounded amount of the response body" do + proxy = Async::HTTP::Proxy.tcp(client, "localhost", 1) + + begin + proxy.connect + rescue Async::HTTP::Proxy::ConnectFailure => error + expect(error.body).to be == reason.byteslice(0, Async::HTTP::Proxy::ConnectFailure::MAXIMUM_BODY_SIZE) + end + + expect(error).to be_a(Async::HTTP::Proxy::ConnectFailure) + + current_task = Async::Task.current + current_task.with_timeout(1) do + current_task.yield while proxy.client.pool.busy? + end + ensure + proxy&.close + end + end + + with "response body which never finishes" do + let(:body) {Async::HTTP::Body::Writable.new} + + let(:app) do + Protocol::HTTP::Middleware.for do |request| + body.write(reason) + + Protocol::HTTP::Response[403, {"content-type" => "text/plain"}, body] + end + end + + it "propagates the caller's timeout" do + proxy = Async::HTTP::Proxy.tcp(client, "localhost", 1) + + expect do + Async::Task.current.with_timeout(0.1) do + proxy.connect + end + end.to raise_exception(Async::TimeoutError) + + current_task = Async::Task.current + current_task.with_timeout(1) do + current_task.yield while proxy.client.pool.busy? + end + ensure + body.close + proxy&.close + end + end + end + + with "data written immediately before closing" do + let(:received) {Async::Promise.new} + + let(:app) do + Protocol::HTTP::Middleware.for do |request| + Async::HTTP::Body::Hijack.response(request, 200, {}) do |stream| + buffer = String.new + + begin + while chunk = stream.read_partial + buffer << chunk + end + ensure + received.resolve(buffer) + end + ensure + stream.close + end + end + end + + def wait_until_released(proxy) + current_task = Async::Task.current + current_task.with_timeout(1) do + received.wait + current_task.yield while proxy.client.pool.busy? + end + end + + it "forwards the data before closing the tunnel" do + proxy = Async::HTTP::Proxy.tcp(client, "localhost", 1) + + peer = proxy.connect + peer.write("QUIT\r\n") + peer.close + + wait_until_released(proxy) + + expect(received.value).to be == "QUIT\r\n" + ensure + proxy&.close + end + + it "forwards the data before closing the tunnel using a block" do + proxy = Async::HTTP::Proxy.tcp(client, "localhost", 1) + + proxy.connect do |peer| + peer.write("QUIT\r\n") + end + + wait_until_released(proxy) + + expect(received.value).to be == "QUIT\r\n" + ensure + proxy&.close + end + end + + with "remote end which closes the tunnel slowly" do + let(:delay) {1} + + let(:app) do + Protocol::HTTP::Middleware.for do |request| + Async::HTTP::Body::Hijack.response(request, 200, {}) do |stream| + stream.read_until("\r\n\r\n") + + stream.write("HTTP/1.1 200 OK\r\ncontent-length: 12\r\n\r\nHello World!") + stream.flush + + # Wait for the client to close its side of the tunnel: + while stream.read_partial + end + + sleep(delay) + ensure + stream.close + end + end + end + + it "closes the tunnel without waiting for the remote end" do + proxied_endpoint = client.proxied_endpoint(Async::HTTP::Endpoint.parse("http://localhost")) + proxied_client = Async::HTTP::Client.new(proxied_endpoint) + + response = proxied_client.get("/") + expect(response.read).to be == "Hello World!" + + proxied_client.close + + start_time = Async::Clock.now + + # The tunnel is released once the request has been forwarded to the proxy: + current_task = Async::Task.current + current_task.with_timeout(delay) do + current_task.yield while client.pool.busy? + end + + client.close + + expect(Async::Clock.now - start_time).to be < (delay / 4.0) + expect(client.pool).to be(:empty?) + end + + it "closes the tunnel without waiting for the remote end using a block" do + proxy = Async::HTTP::Proxy.tcp(client, "localhost", 80) + + proxy.connect do |peer| + peer.write("GET / HTTP/1.1\r\nhost: localhost\r\n\r\n") + + buffer = String.new + buffer << peer.readpartial(1024) until buffer.end_with?("Hello World!") + end + + start_time = Async::Clock.now + + # The tunnel is released once the request has been forwarded to the proxy: + current_task = Async::Task.current + current_task.with_timeout(delay) do + current_task.yield while client.pool.busy? + end + + proxy.close + + expect(Async::Clock.now - start_time).to be < (delay / 4.0) + expect(client.pool).to be(:empty?) + end end with "proxied client" do