From 24a57196c8f74d7f12d9fb9cd43132b5ea8930c1 Mon Sep 17 00:00:00 2001 From: Samuel Williams Date: Thu, 24 Sep 2026 13:09:18 +1200 Subject: [PATCH 01/13] Validate gRPC responses and normalize transport failures --- gems.rb | 3 + lib/async/grpc/client.rb | 37 +++++++++++- lib/async/grpc/transport.rb | 30 ++++++++++ releases.md | 6 ++ test/async/grpc/dispatcher.rb | 2 + test/async/grpc/interoperability.rb | 93 +++++++++++++++++++++++++++++ 6 files changed, 169 insertions(+), 2 deletions(-) create mode 100644 lib/async/grpc/transport.rb create mode 100644 test/async/grpc/interoperability.rb diff --git a/gems.rb b/gems.rb index dc23cba..0312aae 100644 --- a/gems.rb +++ b/gems.rb @@ -7,6 +7,9 @@ gemspec +# Use the protocol fixes while their release is pending: +gem "protocol-grpc", git: "https://github.com/socketry/protocol-grpc.git", ref: "cb474e449da5864a91206a680c2543e84619ea89" + group :maintenance, optional: true do gem "bake-gem" gem "bake-modernize" diff --git a/lib/async/grpc/client.rb b/lib/async/grpc/client.rb index 80f3e98..41f4ba2 100644 --- a/lib/async/grpc/client.rb +++ b/lib/async/grpc/client.rb @@ -17,6 +17,7 @@ require "protocol/grpc/error" require_relative "stub" require_relative "error" +require_relative "transport" module Async module GRPC @@ -101,9 +102,22 @@ def stub(interface_class, service_name) def call(request) request.headers = @headers.merge(request.headers) - super.tap do |response| + response = begin + super + rescue *Transport::ERRORS => error + raise Protocol::GRPC::Unavailable.new(error.message), cause: error + end + + begin response.headers.policy = Protocol::GRPC::HEADER_POLICY + response.body = Transport::Body.new(response.body) if response.body + validate_response!(response) + rescue Exception => error + response.close(error) + raise end + + return response end # Make a gRPC call. @@ -126,7 +140,7 @@ def invoke(service, method, request = nil, metadata: {}, timeout: nil, encoding: headers = Protocol::GRPC::Metadata.build( metadata: metadata, timeout: timeout, - content_type: "application/grpc+proto" + content_type: "application/grpc" ) headers["grpc-encoding"] = encoding if encoding @@ -150,6 +164,25 @@ def invoke(service, method, request = nil, metadata: {}, timeout: nil, encoding: protected + # Reject non-gRPC responses before passing their bytes to a frame decoder. + # @parameter response [Protocol::HTTP::Response] The HTTP response. + # @raises [Protocol::GRPC::Error] If the response is not a valid gRPC envelope. + def validate_response!(response) + content_type = response.headers["content-type"].to_s + return if response.status == 200 && content_type.match?(/\Aapplication\/grpc(?:\+[\w.-]+)?(?:\s*;|\z)/i) + + # Discard raw bytes so trailers remain available without decoding HTML as frames: + response.body&.each{|chunk|} + if response.headers["grpc-status"] + check_status!(response) + status = Protocol::GRPC::Status::INTERNAL + else + status = Protocol::GRPC::Status.for_http_status(response.status) + end + + raise Protocol::GRPC::Error.for(status, "Invalid gRPC response: HTTP #{response.status}, content-type #{content_type.inspect}") + end + # Make a unary gRPC call. # @parameter path [String] The gRPC path # @parameter headers [Protocol::HTTP::Headers] Request headers diff --git a/lib/async/grpc/transport.rb b/lib/async/grpc/transport.rb new file mode 100644 index 0000000..84ca278 --- /dev/null +++ b/lib/async/grpc/transport.rb @@ -0,0 +1,30 @@ +# frozen_string_literal: true + +# Released under the MIT License. +# Copyright, 2026, by Samuel Williams. + +require "socket" +require "openssl" +require "protocol/http/error" +require "protocol/http/body/wrapper" +require "protocol/grpc/error" + +module Async + module GRPC + # Provides error translation at HTTP transport boundaries. + module Transport + ERRORS = [IOError, SystemCallError, SocketError, OpenSSL::SSL::SSLError, Protocol::HTTP::Error].freeze + + # Represents a response body that translates failures reading the transport. + class Body < Protocol::HTTP::Body::Wrapper + # Read raw response bytes, preserving transport failures as the cause. + # @returns [String | Nil] The next chunk, or nil at end of stream. + def read + super + rescue *ERRORS => error + raise Protocol::GRPC::Unavailable.new(error.message), cause: error + end + end + end + end +end diff --git a/releases.md b/releases.md index 61e2c6a..03a4793 100644 --- a/releases.md +++ b/releases.md @@ -1,5 +1,11 @@ # Releases +## Unreleased + + - Send `application/grpc` request headers for compatibility with Google API frontends. + - Validate HTTP responses before decoding gRPC frames, using the HTTP-to-gRPC status fallback when no gRPC status is supplied. + - Translate connection and response-read transport failures into `Protocol::GRPC::Unavailable`, retaining the original exception as the cause. + ## v0.9.0 - Added `Async::GRPC::Dispatcher#emit_completion` for once-per-request completion instrumentation, including routing failures, cancellations, and the final gRPC status. diff --git a/test/async/grpc/dispatcher.rb b/test/async/grpc/dispatcher.rb index 7f7c2f7..20e6272 100644 --- a/test/async/grpc/dispatcher.rb +++ b/test/async/grpc/dispatcher.rb @@ -342,6 +342,7 @@ def close_streams(input, output, call) dispatcher = recording_dispatcher(error_service_name => error_service) response = dispatcher.call(build_request(error_service_name, "RaiseError")) + expect(response.headers["backtrace"]).to be_nil expect(Protocol::GRPC::Metadata.extract_status(response.headers)).to be == Protocol::GRPC::Status::RESOURCE_EXHAUSTED @@ -356,6 +357,7 @@ def close_streams(input, output, call) dispatcher = recording_dispatcher(error_service_name => error_service) response = dispatcher.call(build_request(error_service_name, "RaiseError")) + expect(response.headers["backtrace"]).to be_nil expect(Protocol::GRPC::Metadata.extract_status(response.headers)).to be == Protocol::GRPC::Status::INTERNAL diff --git a/test/async/grpc/interoperability.rb b/test/async/grpc/interoperability.rb new file mode 100644 index 0000000..75f1e75 --- /dev/null +++ b/test/async/grpc/interoperability.rb @@ -0,0 +1,93 @@ +# frozen_string_literal: true + +# Released under the MIT License. +# Copyright, 2026, by Samuel Williams. + +require "async/grpc/client" +require "protocol/http2/error" + +describe Async::GRPC::Client do + let(:request) {Protocol::HTTP::Request["POST", "/example.Service/Call", {}, nil]} + + def client_for(response) + delegate = Object.new + delegate.define_singleton_method(:call){|request| response} + subject.new(delegate) + end + + with "HTTP response validation" do + {400 => 13, 401 => 16, 403 => 7, 404 => 12, 429 => 14, 502 => 14, 503 => 14, 504 => 14, 500 => 2, 200 => 2}.each do |http, grpc| + it "maps HTTP #{http} without decoding an HTML body" do + response = Protocol::HTTP::Response[http, {"content-type" => "text/html"}, [""]] + expect do + client_for(response).call(request) + end.to raise_exception(Protocol::GRPC::Error, message: be =~ /Invalid gRPC response/).and(have_attributes(status_code: be == grpc)) + expect(response.body).to be_nil + end + end + + it "preserves explicit gRPC errors over the HTTP fallback" do + response = Protocol::HTTP::Response[503, {"grpc-status" => "16", "grpc-message" => "Token%20expired"}, []] + expect do + client_for(response).call(request) + end.to raise_exception(Protocol::GRPC::Unauthenticated) + end + + it "reads trailers without decoding a non-gRPC body" do + headers = Protocol::HTTP::Headers.new + body = Protocol::HTTP::Body::Buffered.new(["HTML"]) + body.define_singleton_method(:read) do + chunk = super() + headers["grpc-status"] = "16" unless chunk + chunk + end + response = Protocol::HTTP::Response[503, headers, body] + expect{client_for(response).call(request)}.to raise_exception(Protocol::GRPC::Unauthenticated) + end + + ["application/grpc", "application/grpc+proto", "application/grpc+json", "application/grpc; charset=utf-8"].each do |content_type| + it "accepts #{content_type}" do + response = Protocol::HTTP::Response[200, {"content-type" => content_type, "grpc-status" => "0"}, nil] + expect(client_for(response).call(request)).to be_equal(response) + end + end + + [nil, "application/grpc-web", "text/html"].each do |content_type| + it "rejects invalid content type #{content_type.inspect}" do + headers = {"grpc-status" => "0"} + headers["content-type"] = content_type if content_type + response = Protocol::HTTP::Response[200, headers, []] + expect{client_for(response).call(request)}.to raise_exception(Protocol::GRPC::Internal) + end + end + end + + with "transport errors" do + [Errno::ECONNREFUSED.new, Errno::ECONNRESET.new, EOFError.new, SocketError.new, OpenSSL::SSL::SSLError.new, Protocol::HTTP2::GoawayError.new("Disconnected")].each do |failure| + it "converts #{failure.class} while connecting" do + delegate = Object.new + delegate.define_singleton_method(:call){|request| raise failure} + expect do + subject.new(delegate).call(request) + end.to raise_exception(Protocol::GRPC::Unavailable).and(have_attributes(cause: be_equal(failure))) + end + end + + it "converts connection failures while reading response bytes" do + failure = Errno::ECONNRESET.new + body = Protocol::HTTP::Body::Buffered.new(["unread"]) + body.define_singleton_method(:read){raise failure} + response = Protocol::HTTP::Response[200, {"content-type" => "application/grpc"}, body] + result = client_for(response).call(request) + expect{result.body.read}.to raise_exception(Protocol::GRPC::Unavailable).and(have_attributes(cause: be_equal(failure))) + ensure + result&.close + end + + it "does not convert caller timeouts into unavailable errors" do + delegate = Object.new + delegate.define_singleton_method(:call){|request| raise Async::TimeoutError} + expect{subject.new(delegate).call(request)}.to raise_exception(Async::TimeoutError) + end + end +end From 5b3df0e0cc272ab83698f43a2dbfb6ebddf56cd4 Mon Sep 17 00:00:00 2001 From: Samuel Williams Date: Thu, 24 Sep 2026 13:18:09 +1200 Subject: [PATCH 02/13] End client error messages with exclamation marks --- gems.rb | 2 +- lib/async/grpc/client.rb | 6 +++--- 2 files changed, 4 insertions(+), 4 deletions(-) diff --git a/gems.rb b/gems.rb index 0312aae..8981141 100644 --- a/gems.rb +++ b/gems.rb @@ -8,7 +8,7 @@ gemspec # Use the protocol fixes while their release is pending: -gem "protocol-grpc", git: "https://github.com/socketry/protocol-grpc.git", ref: "cb474e449da5864a91206a680c2543e84619ea89" +gem "protocol-grpc", git: "https://github.com/socketry/protocol-grpc.git", ref: "e0a81da4eb999fc68f9c7e704c363f54230f560a" group :maintenance, optional: true do gem "bake-gem" diff --git a/lib/async/grpc/client.rb b/lib/async/grpc/client.rb index 41f4ba2..4178ef9 100644 --- a/lib/async/grpc/client.rb +++ b/lib/async/grpc/client.rb @@ -134,7 +134,7 @@ def call(request) # @raises [Protocol::GRPC::Error] If the gRPC call fails def invoke(service, method, request = nil, metadata: {}, timeout: nil, encoding: nil, initial: nil, &block) rpc = service.class.lookup_rpc(method) - raise ArgumentError, "Unknown method: #{method}" unless rpc + raise ArgumentError, "Unknown method: #{method}!" unless rpc path = service.path(method) headers = Protocol::GRPC::Metadata.build( @@ -158,7 +158,7 @@ def invoke(service, method, request = nil, metadata: {}, timeout: nil, encoding: when :bidirectional bidirectional_call(path, headers, request_class, response_class, encoding, initial: initial, &block) else - raise ArgumentError, "Unknown streaming type: #{streaming}" + raise ArgumentError, "Unknown streaming type: #{streaming}!" end end @@ -180,7 +180,7 @@ def validate_response!(response) status = Protocol::GRPC::Status.for_http_status(response.status) end - raise Protocol::GRPC::Error.for(status, "Invalid gRPC response: HTTP #{response.status}, content-type #{content_type.inspect}") + raise Protocol::GRPC::Error.for(status, "Invalid gRPC response: HTTP #{response.status}, content-type #{content_type.inspect}!") end # Make a unary gRPC call. From f1f7ff58487e3f79dffd1481ce1665796c6b69db Mon Sep 17 00:00:00 2001 From: Samuel Williams Date: Thu, 24 Sep 2026 13:25:55 +1200 Subject: [PATCH 03/13] Adopt strict metadata decoding and update error assertion --- gems.rb | 2 +- test/async/grpc/client.rb | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/gems.rb b/gems.rb index 8981141..4947270 100644 --- a/gems.rb +++ b/gems.rb @@ -8,7 +8,7 @@ gemspec # Use the protocol fixes while their release is pending: -gem "protocol-grpc", git: "https://github.com/socketry/protocol-grpc.git", ref: "e0a81da4eb999fc68f9c7e704c363f54230f560a" +gem "protocol-grpc", git: "https://github.com/socketry/protocol-grpc.git", ref: "4a35b7cdb9d88e64794d2cf05ba72d9de1140aa9" group :maintenance, optional: true do gem "bake-gem" diff --git a/test/async/grpc/client.rb b/test/async/grpc/client.rb index a7337fa..97b717f 100644 --- a/test/async/grpc/client.rb +++ b/test/async/grpc/client.rb @@ -371,7 +371,7 @@ expect do client.invoke(interface, :InvalidCall) - end.to raise_exception(ArgumentError, message: be == "Unknown streaming type: invalid") + end.to raise_exception(ArgumentError, message: be == "Unknown streaming type: invalid!") end end From 2b31bf4d62917263d9d8f99bcc135079565e056a Mon Sep 17 00:00:00 2001 From: Samuel Williams Date: Thu, 24 Sep 2026 13:38:26 +1200 Subject: [PATCH 04/13] Require the released protocol-grpc 0.16 dependency --- async-grpc.gemspec | 2 +- gems.rb | 3 --- 2 files changed, 1 insertion(+), 4 deletions(-) diff --git a/async-grpc.gemspec b/async-grpc.gemspec index e68dabf..2ae88fd 100644 --- a/async-grpc.gemspec +++ b/async-grpc.gemspec @@ -28,6 +28,6 @@ Gem::Specification.new do |spec| spec.add_dependency "async", ">= 2.38.0" spec.add_dependency "async-http" - spec.add_dependency "protocol-grpc", "~> 0.14" + spec.add_dependency "protocol-grpc", "~> 0.16" spec.add_dependency "protocol-http", "~> 0.60" end diff --git a/gems.rb b/gems.rb index 4947270..dc23cba 100644 --- a/gems.rb +++ b/gems.rb @@ -7,9 +7,6 @@ gemspec -# Use the protocol fixes while their release is pending: -gem "protocol-grpc", git: "https://github.com/socketry/protocol-grpc.git", ref: "4a35b7cdb9d88e64794d2cf05ba72d9de1140aa9" - group :maintenance, optional: true do gem "bake-gem" gem "bake-modernize" From 5ab808b32bf9995b177039c5e15f1a911e2cd6e6 Mon Sep 17 00:00:00 2001 From: Samuel Williams Date: Thu, 24 Sep 2026 13:47:59 +1200 Subject: [PATCH 05/13] Limit client changes to HTTP response validation --- lib/async/grpc/client.rb | 8 +------- lib/async/grpc/transport.rb | 30 ----------------------------- releases.md | 1 - test/async/grpc/interoperability.rb | 30 ----------------------------- 4 files changed, 1 insertion(+), 68 deletions(-) delete mode 100644 lib/async/grpc/transport.rb diff --git a/lib/async/grpc/client.rb b/lib/async/grpc/client.rb index 4178ef9..1b6a506 100644 --- a/lib/async/grpc/client.rb +++ b/lib/async/grpc/client.rb @@ -17,7 +17,6 @@ require "protocol/grpc/error" require_relative "stub" require_relative "error" -require_relative "transport" module Async module GRPC @@ -102,15 +101,10 @@ def stub(interface_class, service_name) def call(request) request.headers = @headers.merge(request.headers) - response = begin - super - rescue *Transport::ERRORS => error - raise Protocol::GRPC::Unavailable.new(error.message), cause: error - end + response = super begin response.headers.policy = Protocol::GRPC::HEADER_POLICY - response.body = Transport::Body.new(response.body) if response.body validate_response!(response) rescue Exception => error response.close(error) diff --git a/lib/async/grpc/transport.rb b/lib/async/grpc/transport.rb deleted file mode 100644 index 84ca278..0000000 --- a/lib/async/grpc/transport.rb +++ /dev/null @@ -1,30 +0,0 @@ -# frozen_string_literal: true - -# Released under the MIT License. -# Copyright, 2026, by Samuel Williams. - -require "socket" -require "openssl" -require "protocol/http/error" -require "protocol/http/body/wrapper" -require "protocol/grpc/error" - -module Async - module GRPC - # Provides error translation at HTTP transport boundaries. - module Transport - ERRORS = [IOError, SystemCallError, SocketError, OpenSSL::SSL::SSLError, Protocol::HTTP::Error].freeze - - # Represents a response body that translates failures reading the transport. - class Body < Protocol::HTTP::Body::Wrapper - # Read raw response bytes, preserving transport failures as the cause. - # @returns [String | Nil] The next chunk, or nil at end of stream. - def read - super - rescue *ERRORS => error - raise Protocol::GRPC::Unavailable.new(error.message), cause: error - end - end - end - end -end diff --git a/releases.md b/releases.md index 03a4793..0923888 100644 --- a/releases.md +++ b/releases.md @@ -4,7 +4,6 @@ - Send `application/grpc` request headers for compatibility with Google API frontends. - Validate HTTP responses before decoding gRPC frames, using the HTTP-to-gRPC status fallback when no gRPC status is supplied. - - Translate connection and response-read transport failures into `Protocol::GRPC::Unavailable`, retaining the original exception as the cause. ## v0.9.0 diff --git a/test/async/grpc/interoperability.rb b/test/async/grpc/interoperability.rb index 75f1e75..7141cb2 100644 --- a/test/async/grpc/interoperability.rb +++ b/test/async/grpc/interoperability.rb @@ -4,7 +4,6 @@ # Copyright, 2026, by Samuel Williams. require "async/grpc/client" -require "protocol/http2/error" describe Async::GRPC::Client do let(:request) {Protocol::HTTP::Request["POST", "/example.Service/Call", {}, nil]} @@ -61,33 +60,4 @@ def client_for(response) end end end - - with "transport errors" do - [Errno::ECONNREFUSED.new, Errno::ECONNRESET.new, EOFError.new, SocketError.new, OpenSSL::SSL::SSLError.new, Protocol::HTTP2::GoawayError.new("Disconnected")].each do |failure| - it "converts #{failure.class} while connecting" do - delegate = Object.new - delegate.define_singleton_method(:call){|request| raise failure} - expect do - subject.new(delegate).call(request) - end.to raise_exception(Protocol::GRPC::Unavailable).and(have_attributes(cause: be_equal(failure))) - end - end - - it "converts connection failures while reading response bytes" do - failure = Errno::ECONNRESET.new - body = Protocol::HTTP::Body::Buffered.new(["unread"]) - body.define_singleton_method(:read){raise failure} - response = Protocol::HTTP::Response[200, {"content-type" => "application/grpc"}, body] - result = client_for(response).call(request) - expect{result.body.read}.to raise_exception(Protocol::GRPC::Unavailable).and(have_attributes(cause: be_equal(failure))) - ensure - result&.close - end - - it "does not convert caller timeouts into unavailable errors" do - delegate = Object.new - delegate.define_singleton_method(:call){|request| raise Async::TimeoutError} - expect{subject.new(delegate).call(request)}.to raise_exception(Async::TimeoutError) - end - end end From dbb55a9fbb96fa7fc7584dc1a8b0ab33845d14d0 Mon Sep 17 00:00:00 2001 From: Samuel Williams Date: Thu, 24 Sep 2026 13:56:04 +1200 Subject: [PATCH 06/13] Use ensure to close responses when validation fails --- lib/async/grpc/client.rb | 12 ++++-------- test/async/grpc/interoperability.rb | 14 +++++++++++++- 2 files changed, 17 insertions(+), 9 deletions(-) diff --git a/lib/async/grpc/client.rb b/lib/async/grpc/client.rb index 1b6a506..0fe4b07 100644 --- a/lib/async/grpc/client.rb +++ b/lib/async/grpc/client.rb @@ -101,17 +101,13 @@ def stub(interface_class, service_name) def call(request) request.headers = @headers.merge(request.headers) - response = super - - begin + super.tap do |response| response.headers.policy = Protocol::GRPC::HEADER_POLICY validate_response!(response) - rescue Exception => error - response.close(error) - raise + success = true + ensure + response.close unless success end - - return response end # Make a gRPC call. diff --git a/test/async/grpc/interoperability.rb b/test/async/grpc/interoperability.rb index 7141cb2..1f9ecea 100644 --- a/test/async/grpc/interoperability.rb +++ b/test/async/grpc/interoperability.rb @@ -32,6 +32,15 @@ def client_for(response) end.to raise_exception(Protocol::GRPC::Unauthenticated) end + it "closes the response when reading an invalid response body fails" do + body = Protocol::HTTP::Body::Buffered.new(["HTML"]) + body.define_singleton_method(:read){raise RuntimeError, "Read failed!"} + response = Protocol::HTTP::Response[503, {"content-type" => "text/html"}, body] + + expect{client_for(response).call(request)}.to raise_exception(RuntimeError, message: be == "Read failed!") + expect(response.body).to be_nil + end + it "reads trailers without decoding a non-gRPC body" do headers = Protocol::HTTP::Headers.new body = Protocol::HTTP::Body::Buffered.new(["HTML"]) @@ -46,8 +55,11 @@ def client_for(response) ["application/grpc", "application/grpc+proto", "application/grpc+json", "application/grpc; charset=utf-8"].each do |content_type| it "accepts #{content_type}" do - response = Protocol::HTTP::Response[200, {"content-type" => content_type, "grpc-status" => "0"}, nil] + response = Protocol::HTTP::Response[200, {"content-type" => content_type, "grpc-status" => "0"}, ["unread"]] expect(client_for(response).call(request)).to be_equal(response) + expect(response.body.read).to be == "unread" + ensure + response.close end end From 76b2e09ea8c6eca7de80c24f0df217f9fc723756 Mon Sep 17 00:00:00 2001 From: Samuel Williams Date: Thu, 24 Sep 2026 14:00:26 +1200 Subject: [PATCH 07/13] Use the body discard API during response validation --- lib/async/grpc/client.rb | 4 ++-- test/async/grpc/interoperability.rb | 8 ++++++-- 2 files changed, 8 insertions(+), 4 deletions(-) diff --git a/lib/async/grpc/client.rb b/lib/async/grpc/client.rb index 0fe4b07..f69d4f4 100644 --- a/lib/async/grpc/client.rb +++ b/lib/async/grpc/client.rb @@ -161,8 +161,8 @@ def validate_response!(response) content_type = response.headers["content-type"].to_s return if response.status == 200 && content_type.match?(/\Aapplication\/grpc(?:\+[\w.-]+)?(?:\s*;|\z)/i) - # Discard raw bytes so trailers remain available without decoding HTML as frames: - response.body&.each{|chunk|} + # Consume the body without decoding frames so gRPC status trailers are available: + response.body&.discard if response.headers["grpc-status"] check_status!(response) status = Protocol::GRPC::Status::INTERNAL diff --git a/test/async/grpc/interoperability.rb b/test/async/grpc/interoperability.rb index 1f9ecea..da6597a 100644 --- a/test/async/grpc/interoperability.rb +++ b/test/async/grpc/interoperability.rb @@ -4,6 +4,7 @@ # Copyright, 2026, by Samuel Williams. require "async/grpc/client" +require "protocol/http/body/writable" describe Async::GRPC::Client do let(:request) {Protocol::HTTP::Request["POST", "/example.Service/Call", {}, nil]} @@ -33,17 +34,20 @@ def client_for(response) end it "closes the response when reading an invalid response body fails" do - body = Protocol::HTTP::Body::Buffered.new(["HTML"]) + body = Protocol::HTTP::Body::Writable.new body.define_singleton_method(:read){raise RuntimeError, "Read failed!"} response = Protocol::HTTP::Response[503, {"content-type" => "text/html"}, body] expect{client_for(response).call(request)}.to raise_exception(RuntimeError, message: be == "Read failed!") expect(response.body).to be_nil + expect(body).to be(:closed?) end it "reads trailers without decoding a non-gRPC body" do headers = Protocol::HTTP::Headers.new - body = Protocol::HTTP::Body::Buffered.new(["HTML"]) + body = Protocol::HTTP::Body::Writable.new + body.write("HTML") + body.close_write body.define_singleton_method(:read) do chunk = super() headers["grpc-status"] = "16" unless chunk From 7a05148ad6cbc5a2a25f7ddbc3d622b18a487c46 Mon Sep 17 00:00:00 2001 From: Samuel Williams Date: Thu, 24 Sep 2026 14:50:13 +1200 Subject: [PATCH 08/13] Preserve invalid HTTP responses in ResponseError --- lib/async/grpc/client.rb | 15 +++--------- lib/async/grpc/error.rb | 21 ++++++++++++++++ releases.md | 2 +- test/async/grpc/interoperability.rb | 38 +++++++++-------------------- 4 files changed, 38 insertions(+), 38 deletions(-) diff --git a/lib/async/grpc/client.rb b/lib/async/grpc/client.rb index f69d4f4..7a96d16 100644 --- a/lib/async/grpc/client.rb +++ b/lib/async/grpc/client.rb @@ -98,6 +98,7 @@ def stub(interface_class, service_name) # Call the underlying HTTP client with merged headers. # @parameter request [Protocol::HTTP::Request] The HTTP request # @returns [Protocol::HTTP::Response] The HTTP response + # @raises [ResponseError] If the HTTP response does not conform to gRPC. def call(request) request.headers = @headers.merge(request.headers) @@ -121,6 +122,7 @@ def call(request) # @yields {|input, output| ...} Block for streaming calls # @returns [Object | Protocol::GRPC::Body::ReadableBody] Response message or readable body for streaming # @raises [ArgumentError] If method is unknown or streaming type is invalid + # @raises [ResponseError] If the HTTP response does not conform to gRPC. # @raises [Protocol::GRPC::Error] If the gRPC call fails def invoke(service, method, request = nil, metadata: {}, timeout: nil, encoding: nil, initial: nil, &block) rpc = service.class.lookup_rpc(method) @@ -156,21 +158,12 @@ def invoke(service, method, request = nil, metadata: {}, timeout: nil, encoding: # Reject non-gRPC responses before passing their bytes to a frame decoder. # @parameter response [Protocol::HTTP::Response] The HTTP response. - # @raises [Protocol::GRPC::Error] If the response is not a valid gRPC envelope. + # @raises [ResponseError] If the response is not a valid gRPC envelope. def validate_response!(response) content_type = response.headers["content-type"].to_s return if response.status == 200 && content_type.match?(/\Aapplication\/grpc(?:\+[\w.-]+)?(?:\s*;|\z)/i) - # Consume the body without decoding frames so gRPC status trailers are available: - response.body&.discard - if response.headers["grpc-status"] - check_status!(response) - status = Protocol::GRPC::Status::INTERNAL - else - status = Protocol::GRPC::Status.for_http_status(response.status) - end - - raise Protocol::GRPC::Error.for(status, "Invalid gRPC response: HTTP #{response.status}, content-type #{content_type.inspect}!") + raise ResponseError, response end # Make a unary gRPC call. diff --git a/lib/async/grpc/error.rb b/lib/async/grpc/error.rb index dfd90e3..18702e1 100644 --- a/lib/async/grpc/error.rb +++ b/lib/async/grpc/error.rb @@ -14,6 +14,27 @@ class Error < StandardError class DeadlineExceededError < Error end + # Raised when an HTTP response does not conform to gRPC. + # Preserves the response body in the error message and the response for inspection. + class ResponseError < Error + # Initialize an error by reading the raw response body. + # @parameter response [Protocol::HTTP::Response] The invalid response. + def initialize(response) + super(response.read.to_s) + + @response = response + end + + # @attribute [Protocol::HTTP::Response] The response, with its body consumed. + attr :response + + # Describe the invalid response and include its body. + # @returns [String] The HTTP status, content type, and response body. + def to_s + "Invalid gRPC response: HTTP #{@response.status}, content-type #{@response.headers["content-type"].to_s.inspect}!\n#{super}" + end + end + # Represents an error that originated from a remote gRPC server. # Used as the `cause` of {Protocol::GRPC::Error} when the client receives a non-OK status. # The message and optional backtrace are extracted from response metadata. diff --git a/releases.md b/releases.md index 0923888..d8a81a6 100644 --- a/releases.md +++ b/releases.md @@ -3,7 +3,7 @@ ## Unreleased - Send `application/grpc` request headers for compatibility with Google API frontends. - - Validate HTTP responses before decoding gRPC frames, using the HTTP-to-gRPC status fallback when no gRPC status is supplied. + - Reject invalid HTTP responses before decoding gRPC frames with `Async::GRPC::ResponseError`, preserving the body in the error message and the response for inspecting status and headers. ## v0.9.0 diff --git a/test/async/grpc/interoperability.rb b/test/async/grpc/interoperability.rb index da6597a..05fa328 100644 --- a/test/async/grpc/interoperability.rb +++ b/test/async/grpc/interoperability.rb @@ -16,21 +16,21 @@ def client_for(response) end with "HTTP response validation" do - {400 => 13, 401 => 16, 403 => 7, 404 => 12, 429 => 14, 502 => 14, 503 => 14, 504 => 14, 500 => 2, 200 => 2}.each do |http, grpc| - it "maps HTTP #{http} without decoding an HTML body" do - response = Protocol::HTTP::Response[http, {"content-type" => "text/html"}, [""]] - expect do - client_for(response).call(request) - end.to raise_exception(Protocol::GRPC::Error, message: be =~ /Invalid gRPC response/).and(have_attributes(status_code: be == grpc)) - expect(response.body).to be_nil - end + it "preserves the HTTP response and HTML body in the error" do + response = Protocol::HTTP::Response[503, {"content-type" => "text/html", "x-request-id" => "123"}, ["", "Proxy failure!", ""]] + expect do + client_for(response).call(request) + end.to raise_exception(Async::GRPC::ResponseError, message: be == "Invalid gRPC response: HTTP 503, content-type \"text/html\"!\nProxy failure!").and(have_attributes(response: be_equal(response))) + expect(response.headers["x-request-id"]).to be == ["123"] + expect(response.body).to be_nil end - it "preserves explicit gRPC errors over the HTTP fallback" do - response = Protocol::HTTP::Response[503, {"grpc-status" => "16", "grpc-message" => "Token%20expired"}, []] + it "rejects a non-200 status even with gRPC content type and OK status" do + response = Protocol::HTTP::Response[503, {"content-type" => "application/grpc", "grpc-status" => "0"}, nil] expect do client_for(response).call(request) - end.to raise_exception(Protocol::GRPC::Unauthenticated) + end.to raise_exception(Async::GRPC::ResponseError, message: be =~ /HTTP 503/) + expect(response.body).to be_nil end it "closes the response when reading an invalid response body fails" do @@ -43,20 +43,6 @@ def client_for(response) expect(body).to be(:closed?) end - it "reads trailers without decoding a non-gRPC body" do - headers = Protocol::HTTP::Headers.new - body = Protocol::HTTP::Body::Writable.new - body.write("HTML") - body.close_write - body.define_singleton_method(:read) do - chunk = super() - headers["grpc-status"] = "16" unless chunk - chunk - end - response = Protocol::HTTP::Response[503, headers, body] - expect{client_for(response).call(request)}.to raise_exception(Protocol::GRPC::Unauthenticated) - end - ["application/grpc", "application/grpc+proto", "application/grpc+json", "application/grpc; charset=utf-8"].each do |content_type| it "accepts #{content_type}" do response = Protocol::HTTP::Response[200, {"content-type" => content_type, "grpc-status" => "0"}, ["unread"]] @@ -72,7 +58,7 @@ def client_for(response) headers = {"grpc-status" => "0"} headers["content-type"] = content_type if content_type response = Protocol::HTTP::Response[200, headers, []] - expect{client_for(response).call(request)}.to raise_exception(Protocol::GRPC::Internal) + expect{client_for(response).call(request)}.to raise_exception(Async::GRPC::ResponseError) end end end From 7de579cbeedc1769e035bc1d997f62b45285366b Mon Sep 17 00:00:00 2001 From: Samuel Williams Date: Thu, 24 Sep 2026 15:00:08 +1200 Subject: [PATCH 09/13] Match async-rest response error initialization --- lib/async/grpc/error.rb | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/lib/async/grpc/error.rb b/lib/async/grpc/error.rb index 18702e1..e6694be 100644 --- a/lib/async/grpc/error.rb +++ b/lib/async/grpc/error.rb @@ -20,7 +20,7 @@ class ResponseError < Error # Initialize an error by reading the raw response body. # @parameter response [Protocol::HTTP::Response] The invalid response. def initialize(response) - super(response.read.to_s) + super(response.read) @response = response end From fc5a2870d3f07db5abbcce4d17886fc40aac26f9 Mon Sep 17 00:00:00 2001 From: Samuel Williams Date: Thu, 24 Sep 2026 15:01:51 +1200 Subject: [PATCH 10/13] Rely on response reads to close invalid response bodies --- lib/async/grpc/client.rb | 3 --- test/async/grpc/interoperability.rb | 3 +-- 2 files changed, 1 insertion(+), 5 deletions(-) diff --git a/lib/async/grpc/client.rb b/lib/async/grpc/client.rb index 7a96d16..5617009 100644 --- a/lib/async/grpc/client.rb +++ b/lib/async/grpc/client.rb @@ -105,9 +105,6 @@ def call(request) super.tap do |response| response.headers.policy = Protocol::GRPC::HEADER_POLICY validate_response!(response) - success = true - ensure - response.close unless success end end diff --git a/test/async/grpc/interoperability.rb b/test/async/grpc/interoperability.rb index 05fa328..32bf7a1 100644 --- a/test/async/grpc/interoperability.rb +++ b/test/async/grpc/interoperability.rb @@ -33,13 +33,12 @@ def client_for(response) expect(response.body).to be_nil end - it "closes the response when reading an invalid response body fails" do + it "closes the body when reading an invalid response fails" do body = Protocol::HTTP::Body::Writable.new body.define_singleton_method(:read){raise RuntimeError, "Read failed!"} response = Protocol::HTTP::Response[503, {"content-type" => "text/html"}, body] expect{client_for(response).call(request)}.to raise_exception(RuntimeError, message: be == "Read failed!") - expect(response.body).to be_nil expect(body).to be(:closed?) end From aae47ed1330463c14193c1740bb9644325faab35 Mon Sep 17 00:00:00 2001 From: Samuel Williams Date: Thu, 24 Sep 2026 15:39:16 +1200 Subject: [PATCH 11/13] Retain buffered responses for error inspection --- lib/async/grpc/error.rb | 16 +++++----------- releases.md | 2 +- test/async/grpc/interoperability.rb | 5 +++-- 3 files changed, 9 insertions(+), 14 deletions(-) diff --git a/lib/async/grpc/error.rb b/lib/async/grpc/error.rb index e6694be..6633be4 100644 --- a/lib/async/grpc/error.rb +++ b/lib/async/grpc/error.rb @@ -15,24 +15,18 @@ class DeadlineExceededError < Error end # Raised when an HTTP response does not conform to gRPC. - # Preserves the response body in the error message and the response for inspection. + # Preserves the buffered response for inspection. class ResponseError < Error - # Initialize an error by reading the raw response body. + # Initialize an error by buffering the raw response body. # @parameter response [Protocol::HTTP::Response] The invalid response. def initialize(response) - super(response.read) + @response = response.buffered! - @response = response + super("Invalid gRPC response: HTTP #{response.status}, content-type #{response.headers["content-type"].to_s.inspect}!") end - # @attribute [Protocol::HTTP::Response] The response, with its body consumed. + # @attribute [Protocol::HTTP::Response] The response, with its body buffered and available to read. attr :response - - # Describe the invalid response and include its body. - # @returns [String] The HTTP status, content type, and response body. - def to_s - "Invalid gRPC response: HTTP #{@response.status}, content-type #{@response.headers["content-type"].to_s.inspect}!\n#{super}" - end end # Represents an error that originated from a remote gRPC server. diff --git a/releases.md b/releases.md index d8a81a6..97d13c8 100644 --- a/releases.md +++ b/releases.md @@ -3,7 +3,7 @@ ## Unreleased - Send `application/grpc` request headers for compatibility with Google API frontends. - - Reject invalid HTTP responses before decoding gRPC frames with `Async::GRPC::ResponseError`, preserving the body in the error message and the response for inspecting status and headers. + - Reject invalid HTTP responses before decoding gRPC frames with `Async::GRPC::ResponseError`, describing the invalid status and content type in the message and buffering the response for inspecting status, headers, and body. ## v0.9.0 diff --git a/test/async/grpc/interoperability.rb b/test/async/grpc/interoperability.rb index 32bf7a1..d4900c3 100644 --- a/test/async/grpc/interoperability.rb +++ b/test/async/grpc/interoperability.rb @@ -20,9 +20,10 @@ def client_for(response) response = Protocol::HTTP::Response[503, {"content-type" => "text/html", "x-request-id" => "123"}, ["", "Proxy failure!", ""]] expect do client_for(response).call(request) - end.to raise_exception(Async::GRPC::ResponseError, message: be == "Invalid gRPC response: HTTP 503, content-type \"text/html\"!\nProxy failure!").and(have_attributes(response: be_equal(response))) + end.to raise_exception(Async::GRPC::ResponseError, message: be == "Invalid gRPC response: HTTP 503, content-type \"text/html\"!").and(have_attributes(response: be_equal(response))) expect(response.headers["x-request-id"]).to be == ["123"] - expect(response.body).to be_nil + expect(response.body).to be_a(Protocol::HTTP::Body::Buffered) + expect(response.read).to be == "Proxy failure!" end it "rejects a non-200 status even with gRPC content type and OK status" do From c93039a9ef3abb9dac36993a9a117a5926c9f7f8 Mon Sep 17 00:00:00 2001 From: Samuel Williams Date: Thu, 24 Sep 2026 15:44:47 +1200 Subject: [PATCH 12/13] Buffer invalid responses at the raise site --- lib/async/grpc/client.rb | 2 +- lib/async/grpc/error.rb | 6 +++--- 2 files changed, 4 insertions(+), 4 deletions(-) diff --git a/lib/async/grpc/client.rb b/lib/async/grpc/client.rb index 5617009..0056f15 100644 --- a/lib/async/grpc/client.rb +++ b/lib/async/grpc/client.rb @@ -160,7 +160,7 @@ def validate_response!(response) content_type = response.headers["content-type"].to_s return if response.status == 200 && content_type.match?(/\Aapplication\/grpc(?:\+[\w.-]+)?(?:\s*;|\z)/i) - raise ResponseError, response + raise ResponseError, response.buffered! end # Make a unary gRPC call. diff --git a/lib/async/grpc/error.rb b/lib/async/grpc/error.rb index 6633be4..dc98e34 100644 --- a/lib/async/grpc/error.rb +++ b/lib/async/grpc/error.rb @@ -17,10 +17,10 @@ class DeadlineExceededError < Error # Raised when an HTTP response does not conform to gRPC. # Preserves the buffered response for inspection. class ResponseError < Error - # Initialize an error by buffering the raw response body. - # @parameter response [Protocol::HTTP::Response] The invalid response. + # Initialize an error from an invalid response. + # @parameter response [Protocol::HTTP::Response] The invalid response, with its body already buffered. def initialize(response) - @response = response.buffered! + @response = response super("Invalid gRPC response: HTTP #{response.status}, content-type #{response.headers["content-type"].to_s.inspect}!") end From 90623d0037aa7d5ae109cc79f72a57620e663577 Mon Sep 17 00:00:00 2001 From: Samuel Williams Date: Thu, 24 Sep 2026 15:46:03 +1200 Subject: [PATCH 13/13] Build response errors with a factory and explicit message --- lib/async/grpc/client.rb | 2 +- lib/async/grpc/error.rb | 20 ++++++++++++++------ 2 files changed, 15 insertions(+), 7 deletions(-) diff --git a/lib/async/grpc/client.rb b/lib/async/grpc/client.rb index 0056f15..eb831ed 100644 --- a/lib/async/grpc/client.rb +++ b/lib/async/grpc/client.rb @@ -160,7 +160,7 @@ def validate_response!(response) content_type = response.headers["content-type"].to_s return if response.status == 200 && content_type.match?(/\Aapplication\/grpc(?:\+[\w.-]+)?(?:\s*;|\z)/i) - raise ResponseError, response.buffered! + raise ResponseError.for(response) end # Make a unary gRPC call. diff --git a/lib/async/grpc/error.rb b/lib/async/grpc/error.rb index dc98e34..1c76ec8 100644 --- a/lib/async/grpc/error.rb +++ b/lib/async/grpc/error.rb @@ -17,15 +17,23 @@ class DeadlineExceededError < Error # Raised when an HTTP response does not conform to gRPC. # Preserves the buffered response for inspection. class ResponseError < Error - # Initialize an error from an invalid response. - # @parameter response [Protocol::HTTP::Response] The invalid response, with its body already buffered. - def initialize(response) - @response = response + # Create an error from an invalid response, buffering its body for inspection. + # @parameter response [Protocol::HTTP::Response] The invalid response. + # @returns [ResponseError] The error with the buffered response attached. + def self.for(response) + self.new("Invalid gRPC response: HTTP #{response.status}, content-type #{response.headers["content-type"].to_s.inspect}!", response.buffered!) + end + + # Initialize an error with a message and response. + # @parameter message [String] The reason the response is invalid. + # @parameter response [Protocol::HTTP::Response] The buffered response. + def initialize(message, response) + super(message) - super("Invalid gRPC response: HTTP #{response.status}, content-type #{response.headers["content-type"].to_s.inspect}!") + @response = response end - # @attribute [Protocol::HTTP::Response] The response, with its body buffered and available to read. + # @attribute [Protocol::HTTP::Response] The buffered response, with its body available to read. attr :response end