From 554ccc9a564e824ac4c18a042583f449e17dd73f Mon Sep 17 00:00:00 2001 From: Martin Ek Date: Wed, 30 Sep 2026 00:20:58 -0700 Subject: [PATCH 1/2] Close the CONNECT tunnel when a proxied socket is closed. Closing the socket returned by `Proxy#connect` without a block only half-closed the CONNECT request body, so the proxy client's connection stayed leased (and `Client#close` blocked) until the proxy closed its end of the tunnel. Closing it now stops the pipe and closes the CONNECT request and response with an error. The block form closes the same way. `ConnectFailure#body` also exposes up to 8 KiB of the proxy's response body. Co-Authored-By: Claude Opus 5.5 (1M context) --- lib/async/http/body/pipe.rb | 11 ++- lib/async/http/proxy.rb | 82 +++++++++++++++- releases.md | 5 + test/async/http/body/pipe.rb | 39 ++++++++ test/async/http/proxy.rb | 184 ++++++++++++++++++++++++++++++++++- 5 files changed, 310 insertions(+), 11 deletions(-) diff --git a/lib/async/http/body/pipe.rb b/lib/async/http/body/pipe.rb index da221778..b17291be 100644 --- a/lib/async/http/body/pipe.rb +++ b/lib/async/http/body/pipe.rb @@ -24,6 +24,8 @@ def initialize(input, output = Writable.new, task: Task.current) @reader = nil @writer = nil + @error = nil + task.async(transient: true, &self.method(:reader)) task.async(transient: true, &self.method(:writer)) end @@ -34,7 +36,10 @@ def to_io end # Close the pipe and stop the reader and writer tasks. - def close + # @parameter error [Exception | Nil] If given, the input and output bodies are closed with this error, rather than being closed normally. + def close(error = nil) + @error ||= error + @reader&.stop @writer&.stop @@ -57,7 +62,7 @@ def reader(task) @head.close_write rescue => error ensure - @input.close(error) + @input.close(error || @error) close_head if @writer&.finished? end @@ -74,7 +79,7 @@ def writer(task) end rescue => error ensure - @output.close_write(error) + @output.close_write(error || @error) close_head if @reader&.finished? end diff --git a/lib/async/http/proxy.rb b/lib/async/http/proxy.rb index 4f295f7e..c45fda1f 100644 --- a/lib/async/http/proxy.rb +++ b/lib/async/http/proxy.rb @@ -15,14 +15,60 @@ 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 closes the pipe, and the CONNECT request and response with an error, so that HTTP/2 resets the stream and HTTP/1 closes the connection, neither of which can be reused after a tunnel anyway. + module Tunnel + # Used to close the CONNECT request and response when the socket is closed. + class Closed < IOError + end + + # Close the pipe, and the CONNECT request and response with an error. + # @parameter pipe [Body::Pipe] The pipe connected to the CONNECT request and response. + def self.close(pipe) + pipe.close(Closed.new("Tunnel closed!")) + end + + # @parameter pipe [Body::Pipe] The pipe connected to the CONNECT request and response. + # @returns [Socket] The socket for the pipe, which closes the pipe when closed. + def self.wrap(pipe) + socket = pipe.to_io + socket.extend(self) + socket.instance_variable_set(:@pipe, pipe) + + return socket + end + + # Close the pipe, which also closes the socket. + def close + if pipe = @pipe + @pipe = nil + Tunnel.close(pipe) + else + super + end + end end # Extends {Async::HTTP::Client} with proxy capabilities. @@ -99,19 +145,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 + Tunnel.close(pipe) 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 +170,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..d8257e48 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 with an error, so HTTP/2 proxies receive `RST_STREAM` and HTTP/1 proxy connections are closed, and `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. + - `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..ec2cb26d 100644 --- a/test/async/http/body/pipe.rb +++ b/test/async/http/body/pipe.rb @@ -54,6 +54,45 @@ def before 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 and output normally" do + pipe.close + + expect(input.read).to be_nil + expect(output.read).to be_nil + end + + it "closes the input and output with the given error" do + pipe.close(IOError.new("Pipe closed!")) + + expect{input.read}.to raise_exception(IOError, message: be == "Pipe closed!") + expect{output.read}.to raise_exception(IOError, message: be == "Pipe closed!") + end + + 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..26fe3346 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,180 @@ 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 + write_response.resolve(true) unless write_response.resolved? + 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 "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 + 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 + + # HTTP/1 releases the connection once the tunnel's writer task observes the closed request body: + 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 @@ -206,6 +381,9 @@ upstream.write(chunk) upstream.flush end + rescue ::Protocol::HTTP2::StreamError => error + # The client resets the tunnel when it closes it: + Console.debug(self, "Tunnel reset by client.", error) ensure Console.debug(self){"Finished writing to upstream..."} upstream.close_write unless upstream.closed? From 351ef6d433abaf0b95d086a06bef2dcb895e0bd3 Mon Sep 17 00:00:00 2001 From: Martin Ek Date: Wed, 30 Sep 2026 08:44:46 -0700 Subject: [PATCH 2/2] Forward pending writes before closing a proxy tunnel. Closing the tunnel stopped the pipe's writer and closed the CONNECT request with an error, discarding any data written to the socket that hadn't been forwarded yet (e.g. `QUIT\r\n` followed immediately by `close`). Instead, stop only the reader and close the socket, so the writer forwards the remaining data and ends the request normally. HTTP/2 then resets the stream with `NO_ERROR` and HTTP/1 closes the connection. The block form closes the same way. Co-Authored-By: Claude Opus 5.5 (1M context) --- lib/async/http/body/pipe.rb | 18 +++++---- lib/async/http/proxy.rb | 20 +++------- releases.md | 2 +- test/async/http/body/pipe.rb | 24 +++++++----- test/async/http/proxy.rb | 71 +++++++++++++++++++++++++++++++++--- 5 files changed, 96 insertions(+), 39 deletions(-) diff --git a/lib/async/http/body/pipe.rb b/lib/async/http/body/pipe.rb index b17291be..158cdb27 100644 --- a/lib/async/http/body/pipe.rb +++ b/lib/async/http/body/pipe.rb @@ -24,8 +24,6 @@ def initialize(input, output = Writable.new, task: Task.current) @reader = nil @writer = nil - @error = nil - task.async(transient: true, &self.method(:reader)) task.async(transient: true, &self.method(:writer)) end @@ -35,11 +33,15 @@ def to_io @tail end - # Close the pipe and stop the reader and writer tasks. - # @parameter error [Exception | Nil] If given, the input and output bodies are closed with this error, rather than being closed normally. - def close(error = nil) - @error ||= error + # 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 @writer&.stop @@ -62,7 +64,7 @@ def reader(task) @head.close_write rescue => error ensure - @input.close(error || @error) + @input.close(error) close_head if @writer&.finished? end @@ -79,7 +81,7 @@ def writer(task) end rescue => error ensure - @output.close_write(error || @error) + @output.close_write(error) close_head if @reader&.finished? end diff --git a/lib/async/http/proxy.rb b/lib/async/http/proxy.rb index c45fda1f..f42151fe 100644 --- a/lib/async/http/proxy.rb +++ b/lib/async/http/proxy.rb @@ -38,20 +38,10 @@ def initialize(response, body = nil) # 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 closes the pipe, and the CONNECT request and response with an error, so that HTTP/2 resets the stream and HTTP/1 closes the connection, neither of which can be reused after a tunnel anyway. + # 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 - # Used to close the CONNECT request and response when the socket is closed. - class Closed < IOError - end - - # Close the pipe, and the CONNECT request and response with an error. - # @parameter pipe [Body::Pipe] The pipe connected to the CONNECT request and response. - def self.close(pipe) - pipe.close(Closed.new("Tunnel closed!")) - end - # @parameter pipe [Body::Pipe] The pipe connected to the CONNECT request and response. - # @returns [Socket] The socket for the pipe, which closes the pipe when closed. + # @returns [Socket] The socket for the pipe, which finishes the pipe when closed. def self.wrap(pipe) socket = pipe.to_io socket.extend(self) @@ -60,11 +50,11 @@ def self.wrap(pipe) return socket end - # Close the pipe, which also closes the socket. + # Finish the pipe, which also closes the socket. def close if pipe = @pipe @pipe = nil - Tunnel.close(pipe) + pipe.finish else super end @@ -150,7 +140,7 @@ def connect(&block) begin yield pipe.to_io ensure - Tunnel.close(pipe) + pipe.finish end else # This ensures we don't leave a response dangling: diff --git a/releases.md b/releases.md index d8257e48..76e01a2c 100644 --- a/releases.md +++ b/releases.md @@ -2,7 +2,7 @@ ## 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 with an error, so HTTP/2 proxies receive `RST_STREAM` and HTTP/1 proxy connections are closed, and `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. + - 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 diff --git a/test/async/http/body/pipe.rb b/test/async/http/body/pipe.rb index ec2cb26d..1542e394 100644 --- a/test/async/http/body/pipe.rb +++ b/test/async/http/body/pipe.rb @@ -54,25 +54,29 @@ def before end end - with "#close" do + with "#finish" do include Sus::Fixtures::Async::ReactorContext let(:output) {Async::HTTP::Body::Writable.new} let(:pipe) {subject.new(input, output)} - it "closes the input and output normally" do - pipe.close + 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(input.read).to be_nil + expect(output.read).to be == data expect(output.read).to be_nil end + end + + with "#close" do + include Sus::Fixtures::Async::ReactorContext - it "closes the input and output with the given error" do - pipe.close(IOError.new("Pipe closed!")) - - expect{input.read}.to raise_exception(IOError, message: be == "Pipe closed!") - expect{output.read}.to raise_exception(IOError, message: be == "Pipe closed!") - end + 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 diff --git a/test/async/http/proxy.rb b/test/async/http/proxy.rb index 26fe3346..2f9c8242 100644 --- a/test/async/http/proxy.rb +++ b/test/async/http/proxy.rb @@ -188,7 +188,6 @@ expect(write_response).not.to be(:resolved?) expect(proxy.client.pool).not.to be(:busy?) ensure - write_response.resolve(true) unless write_response.resolved? peer&.close proxy&.close end @@ -281,6 +280,64 @@ 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} @@ -313,6 +370,13 @@ 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) @@ -331,7 +395,7 @@ start_time = Async::Clock.now - # HTTP/1 releases the connection once the tunnel's writer task observes the closed request body: + # 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? @@ -381,9 +445,6 @@ upstream.write(chunk) upstream.flush end - rescue ::Protocol::HTTP2::StreamError => error - # The client resets the tunnel when it closes it: - Console.debug(self, "Tunnel reset by client.", error) ensure Console.debug(self){"Finished writing to upstream..."} upstream.close_write unless upstream.closed?