Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 7 additions & 0 deletions lib/async/http/body/pipe.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
72 changes: 67 additions & 5 deletions lib/async/http/proxy.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -99,26 +135,52 @@ 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

# @return [Async::HTTP::Endpoint] an endpoint that connects via the specified proxy.
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)
Expand Down
5 changes: 5 additions & 0 deletions releases.md
Original file line number Diff line number Diff line change
@@ -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.
Expand Down
43 changes: 43 additions & 0 deletions test/async/http/body/pipe.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading