From d47c8dfb926e2239e4711f24c3f2b6bdc831ebe6 Mon Sep 17 00:00:00 2001 From: Martin Ek Date: Thu, 24 Sep 2026 10:36:55 -0700 Subject: [PATCH 1/4] Map transport failures to grpc-ruby errors. Connection, DNS, TLS, and HTTP/2 connection failures now raise `GRPC::Unavailable`, and HTTP/2 stream resets use gRPC's HTTP/2 status mapping. The original exception is preserved as the cause, so retry policies such as GAPIC's `UNAVAILABLE` retries apply. The gRPC response body treats a failed stream as empty, so wrap it to ensure stream errors are raised rather than reported as a missing status. Fixes #5. Co-Authored-By: Claude Opus 5.5 (1M context) --- lib/async/grpc/compatible/client_stub.rb | 48 ++++++- test/async/grpc/compatible/client_stub.rb | 146 +++++++++++++++++++++- 2 files changed, 190 insertions(+), 4 deletions(-) diff --git a/lib/async/grpc/compatible/client_stub.rb b/lib/async/grpc/compatible/client_stub.rb index 8862bc7..0570648 100644 --- a/lib/async/grpc/compatible/client_stub.rb +++ b/lib/async/grpc/compatible/client_stub.rb @@ -7,6 +7,7 @@ require "async/http/endpoint" require "async/http/protocol/http2" require "base64" +require "openssl" require "grpc" require "io/endpoint/tls/configuration" require "protocol/grpc/body/readable" @@ -49,6 +50,35 @@ class ClientStub INSECURE_CREDENTIALS = :this_channel_is_insecure DEFAULT_TIMEOUT = nil + # Transport failures which end a call without a gRPC status. grpc-ruby reports these as `UNAVAILABLE`. + TRANSPORT_ERRORS = [ + ::Protocol::HTTP::RefusedError, + ::Protocol::HTTP2::Error, + ::Protocol::HPACK::Error, + IOError, + SocketError, + SystemCallError, + OpenSSL::SSL::SSLError, + ].freeze + + # The gRPC status for each HTTP/2 stream reset code, from gRPC's HTTP/2 status mapping. Other codes map to `INTERNAL`. + STREAM_RESET_STATUSES = { + ::Protocol::HTTP2::Error::REFUSED_STREAM => ::GRPC::Core::StatusCodes::UNAVAILABLE, + ::Protocol::HTTP2::Error::CANCEL => ::GRPC::Core::StatusCodes::CANCELLED, + ::Protocol::HTTP2::Error::ENHANCE_YOUR_CALM => ::GRPC::Core::StatusCodes::RESOURCE_EXHAUSTED, + ::Protocol::HTTP2::Error::INADEQUATE_SECURITY => ::GRPC::Core::StatusCodes::PERMISSION_DENIED, + }.freeze + + # Represents a response body which is never reported as empty. The gRPC body skips reading an empty body, which would hide the error from a failed stream. + class ResponseBody < ::Protocol::HTTP::Body::Wrapper + # @returns [Boolean] Always false, so that reads continue until the stream ends or fails. + def empty? + false + end + end + + private_constant :ResponseBody + # Build a compatible stub class for a generated GRPC::GenericService. # @parameter service [Class] The generated service definition. # @returns [Class] A client stub with methods for the service's unary RPCs. @@ -285,11 +315,19 @@ def invoke_request_response(method, request, marshal, unmarshal, metadata, timeo content_type: "application/grpc" ) request = Protocol::HTTP::Request["POST", normalize_method(method), headers, body] + payload = perform_request(request, operation) + + payload ? unmarshal.call(payload) : nil + end + + # Send the request and read its response payload. Application callbacks run outside this method, so their errors are not treated as transport failures. + def perform_request(request, operation) response = @channel.client.call(request) begin operation.metadata = extract_metadata(Protocol::HTTP::Headers.new(response.headers.header.to_a, policy: Protocol::GRPC::HEADER_POLICY)) response_encoding = response.headers["grpc-encoding"] + response.body = ResponseBody.new(response.body) if response.body response_body = Protocol::GRPC::Body::Readable.wrap(response, encoding: response_encoding) payload = response_body&.read response_body&.finish @@ -302,10 +340,18 @@ def invoke_request_response(method, request, marshal, unmarshal, metadata, timeo ) check_status!(response) - payload ? unmarshal.call(payload) : nil + payload ensure response.close end + rescue ::Protocol::HTTP2::StreamError => error + status = STREAM_RESET_STATUSES.fetch(error.code, ::GRPC::Core::StatusCodes::INTERNAL) + raise_bad_status(status, error.message, {}, cause: error) + rescue ::Protocol::HTTP::RemoteError => error + # The peer reset the stream with `INTERNAL_ERROR`: + raise_bad_status(::GRPC::Core::StatusCodes::INTERNAL, error.message, {}, cause: error) + rescue *TRANSPORT_ERRORS => error + raise_bad_status(::GRPC::Core::StatusCodes::UNAVAILABLE, error.message, {}, cause: error) end def check_status!(response) diff --git a/test/async/grpc/compatible/client_stub.rb b/test/async/grpc/compatible/client_stub.rb index 432548a..311d970 100644 --- a/test/async/grpc/compatible/client_stub.rb +++ b/test/async/grpc/compatible/client_stub.rb @@ -448,6 +448,121 @@ def request(value, **options) end end + with "transport failures" do + def request_with(client) + subject.new("unused", nil, channel_override: Async::GRPC::Compatible::Channel.new(client: client)).request_response( + "/#{service_name}/Echo", + CompatibleMessage.new("Hello"), + CompatibleMessage.method(:encode), + CompatibleMessage.method(:decode), + return_op: true + ) + end + + def request_failing_with(error) + failing_client = Object.new + failing_client.define_singleton_method(:call){|_request| raise error} + request_with(failing_client) + end + + [ + [Errno::ECONNREFUSED.new("connect(2) for 127.0.0.1:1"), ::GRPC::Unavailable], + [Errno::ECONNRESET.new, ::GRPC::Unavailable], + [Errno::EPIPE.new, ::GRPC::Unavailable], + [EOFError.new("Stream closed before response headers were received!"), ::GRPC::Unavailable], + [IOError.new("Connection closed!"), ::GRPC::Unavailable], + [SocketError.new("getaddrinfo: Name or service not known"), ::GRPC::Unavailable], + [OpenSSL::SSL::SSLError.new("certificate verify failed"), ::GRPC::Unavailable], + [Protocol::HTTP::RefusedError.new("GOAWAY: request not processed."), ::GRPC::Unavailable], + [Protocol::HTTP2::GoawayError.new("Shutting down!", Protocol::HTTP2::Error::INTERNAL_ERROR), ::GRPC::Unavailable], + [Protocol::HTTP2::ProtocolError.new("Invalid frame!"), ::GRPC::Unavailable], + [Protocol::HPACK::Error.new("Invalid index!"), ::GRPC::Unavailable], + [Protocol::HTTP2::StreamError.for(Protocol::HTTP2::Error::REFUSED_STREAM), ::GRPC::Unavailable], + [Protocol::HTTP2::StreamError.for(Protocol::HTTP2::Error::CANCEL), ::GRPC::Cancelled], + [Protocol::HTTP2::StreamError.for(Protocol::HTTP2::Error::ENHANCE_YOUR_CALM), ::GRPC::ResourceExhausted], + [Protocol::HTTP2::StreamError.for(Protocol::HTTP2::Error::INADEQUATE_SECURITY), ::GRPC::PermissionDenied], + [Protocol::HTTP2::StreamError.for(Protocol::HTTP2::Error::PROTOCOL_ERROR), ::GRPC::Internal], + [Protocol::HTTP::RemoteError.new("Internal error!"), ::GRPC::Internal], + ].each do |error, error_class| + it "maps #{error.class}: #{error.message} to #{error_class}" do + operation = request_failing_with(error) + + expect{operation.execute}.to raise_exception(error_class, message: be(:include?, error.message)).and(have_attributes( + cause: be(:equal?, error) + )) + expect(operation.status.code).to be == error_class.new.code + end + end + + it "raises unavailable when the connection fails while reading the response" do + error = Errno::ECONNRESET.new + body = Protocol::HTTP::Body::Writable.new + body.close_write(error) + client = Object.new + client.define_singleton_method(:call){|_request| Protocol::HTTP::Response[200, {"content-type" => "application/grpc"}, body]} + + expect{request_with(client).execute}.to raise_exception(::GRPC::Unavailable).and(have_attributes( + cause: be(:equal?, error) + )) + end + + it "raises unavailable when the connection is refused" do + server = TCPServer.new("127.0.0.1", 0) + port = server.local_address.ip_port + server.close + + direct_stub = subject.new("127.0.0.1:#{port}", :this_channel_is_insecure) + + expect do + direct_stub.request_response("/#{service_name}/Echo", CompatibleMessage.new("Hello"), CompatibleMessage.method(:encode), CompatibleMessage.method(:decode)) + end.to raise_exception(::GRPC::Unavailable, message: be =~ /Connection refused/).and(have_attributes( + cause: be_a(Errno::ECONNREFUSED) + )) + ensure + direct_stub&.close + end + + it "raises unavailable when the server closes the connection" do + server = TCPServer.new("127.0.0.1", 0) + acceptor = Async{server.accept.close} + direct_stub = subject.new("127.0.0.1:#{server.local_address.ip_port}", :this_channel_is_insecure) + + expect do + direct_stub.request_response("/#{service_name}/Echo", CompatibleMessage.new("Hello"), CompatibleMessage.method(:encode), CompatibleMessage.method(:decode)) + end.to raise_exception(::GRPC::Unavailable).and(have_attributes( + cause: be_a(IOError).or(be_a(SystemCallError)) + )) + ensure + direct_stub&.close + acceptor&.stop + server&.close + end + + it "preserves application encoder errors" do + expect do + stub.request_response("/#{service_name}/Echo", CompatibleMessage.new("Hello"), ->(message){raise IOError, "Application encoder!"}, CompatibleMessage.method(:decode)) + end.to raise_exception(IOError, message: be == "Application encoder!") + end + end + + with "a server which resets the stream after sending response headers" do + let(:app) do + Protocol::HTTP::Middleware.for do |request| + body = Protocol::HTTP::Body::Writable.new + body.write("\x00") + body.close_write(RuntimeError.new("Server failure!")) + + Protocol::HTTP::Response[200, {"content-type" => "application/grpc"}, body] + end + end + + it "maps the stream reset to an internal error" do + expect{request("Hello")}.to raise_exception(::GRPC::Internal).and(have_attributes( + cause: be_a(Protocol::HTTP2::StreamError).and(have_attributes(code: be == Protocol::HTTP2::Error::INTERNAL_ERROR)) + )) + end + end + it "defers execution until the operation executes" do operation = request("Hello", return_op: true) expect(operation.status).to be_nil @@ -525,6 +640,25 @@ def request(value, **options) gapic.close end + it "retries transport failures using the GAPIC retry policy" do + attempts = 0 + mock(grpc_client) do |wrapper| + wrapper.wrap(:call) do |original, request| + attempts += 1 + raise Errno::ECONNRESET if attempts == 1 + original.call(request) + end + end + retry_policy = {retry_codes: [::GRPC::Core::StatusCodes::UNAVAILABLE], initial_delay: 0.001, max_delay: 0.001} + + response = gapic.call_rpc(:echo, CompatibleMessage.new("retried"), options: {retry_policy: retry_policy}) + + expect(response.value).to be == "retried" + expect(attempts).to be == 2 + ensure + gapic.close + end + it "supports the generated service helper directly" do generated = subject.for(GeneratedCompatibleService).new("unused", nil, channel_override: channel) expect(generated.echo(CompatibleMessage.new("generated")).value).to be == "generated" @@ -588,7 +722,9 @@ def request(value, **options) expect do direct_stub.request_response("/#{service_name}/Echo", CompatibleMessage.new("TLS"), CompatibleMessage.method(:encode), CompatibleMessage.method(:decode)) - end.to raise_exception(OpenSSL::SSL::SSLError, message: be =~ /certificate verify failed/) + end.to raise_exception(::GRPC::Unavailable, message: be =~ /certificate verify failed/).and(have_attributes( + cause: be_a(OpenSSL::SSL::SSLError) + )) ensure direct_stub&.close end @@ -601,7 +737,9 @@ def request(value, **options) expect do wrong_stub.request_response("/#{service_name}/Echo", CompatibleMessage.new("TLS"), CompatibleMessage.method(:encode), CompatibleMessage.method(:decode)) - end.to raise_exception(OpenSSL::SSL::SSLError, message: be =~ /hostname mismatch/) + end.to raise_exception(::GRPC::Unavailable, message: be =~ /hostname mismatch/).and(have_attributes( + cause: be_a(OpenSSL::SSL::SSLError) + )) ensure wrong_channel&.close end @@ -644,7 +782,9 @@ def server_tls_configuration expect do direct_stub.request_response("/#{service_name}/Echo", CompatibleMessage.new("mTLS"), CompatibleMessage.method(:encode), CompatibleMessage.method(:decode)) - end.to raise_exception(StandardError).and(be_a(OpenSSL::SSL::SSLError).or(be_a(EOFError))) + end.to raise_exception(::GRPC::Unavailable).and(have_attributes( + cause: be_a(OpenSSL::SSL::SSLError).or(be_a(IOError), be_a(SystemCallError)) + )) ensure direct_stub&.close end From c66969377ef8280fe65c11e399d788362922c332 Mon Sep 17 00:00:00 2001 From: Martin Ek Date: Thu, 24 Sep 2026 10:38:37 -0700 Subject: [PATCH 2/4] Remove outdated transport error note from readme. Co-Authored-By: Claude Opus 5.5 (1M context) --- readme.md | 2 -- 1 file changed, 2 deletions(-) diff --git a/readme.md b/readme.md index ceb5cc5..53d236f 100644 --- a/readme.md +++ b/readme.md @@ -60,8 +60,6 @@ The following are not yet supported: Invalid HTTP responses become `GRPC::BadStatus` subclasses using the HTTP status mapping. The error details describe the invalid HTTP status and content type, and `error.cause` is an `Async::GRPC::ResponseError` whose `response` exposes the HTTP status, headers, and buffered body. Call `error.cause.response.read` to read that body. -Socket and TLS failures can still raise native Ruby exceptions. Translation into grpc-ruby transport errors is tracked separately in [issue #5](https://github.com/socketry/async-grpc-compatible/issues/5). - ## Operations and credentials Pass `return_op: true` to defer a unary call until `operation.execute`. An operation executes once and exposes `deadline`, `metadata`, `trailing_metadata`, `status`, `cancel`, and `cancelled?`. The deadline includes time spent waiting to execute. Cancel an active operation from the same Async reactor; cancelling it closes that call without closing a shared channel. Calling `cancel` after completion has no effect. From b92fd9de8fdc0c69fa88b6aa11a49b6453b28f75 Mon Sep 17 00:00:00 2001 From: Martin Ek Date: Thu, 24 Sep 2026 10:48:40 -0700 Subject: [PATCH 3/4] Add release note for transport error mapping. Co-Authored-By: Claude Opus 5.5 (1M context) --- releases.md | 1 + 1 file changed, 1 insertion(+) diff --git a/releases.md b/releases.md index 2f6da39..1600352 100644 --- a/releases.md +++ b/releases.md @@ -9,6 +9,7 @@ - Support custom trust roots and mutual TLS through `IO::Endpoint::TLS::Configuration`. Reject opaque native credentials, conflicting target schemes, and authentication callbacks on plaintext channels. - Map grpc-ruby's TLS constructor arguments with `Compatible::ChannelCredentials.new(root_certificates, private_key, certificate_chain)`, returning an `IO::Endpoint::TLS::Configuration` with custom roots, client certificate chains, and peer verification enabled. - Add `ClientStub.for(service)` and the optional `GapicServiceStub` adapter for generated services and GAPIC clients. + - Map transport failures to grpc-ruby errors, preserving the original exception as the cause. Connection, DNS, TLS, and HTTP/2 connection failures become `GRPC::Unavailable`, and HTTP/2 stream resets use gRPC's HTTP/2 status mapping. ## v0.0.0 From 681eac01aa05a683f6b5af10497b224bbd2aa160 Mon Sep 17 00:00:00 2001 From: Samuel Williams Date: Fri, 25 Sep 2026 10:27:27 +1200 Subject: [PATCH 4/4] Use protocol-grpc body error propagation --- async-grpc-compatible.gemspec | 1 + lib/async/grpc/compatible/client_stub.rb | 11 ----------- 2 files changed, 1 insertion(+), 11 deletions(-) diff --git a/async-grpc-compatible.gemspec b/async-grpc-compatible.gemspec index 1f33bd9..93bc714 100644 --- a/async-grpc-compatible.gemspec +++ b/async-grpc-compatible.gemspec @@ -25,4 +25,5 @@ Gem::Specification.new do |specification| specification.add_dependency "async-http", "~> 0.100" specification.add_dependency "grpc" specification.add_dependency "io-endpoint", "~> 0.18" + specification.add_dependency "protocol-grpc", "~> 0.17" end diff --git a/lib/async/grpc/compatible/client_stub.rb b/lib/async/grpc/compatible/client_stub.rb index 0570648..dfefadc 100644 --- a/lib/async/grpc/compatible/client_stub.rb +++ b/lib/async/grpc/compatible/client_stub.rb @@ -69,16 +69,6 @@ class ClientStub ::Protocol::HTTP2::Error::INADEQUATE_SECURITY => ::GRPC::Core::StatusCodes::PERMISSION_DENIED, }.freeze - # Represents a response body which is never reported as empty. The gRPC body skips reading an empty body, which would hide the error from a failed stream. - class ResponseBody < ::Protocol::HTTP::Body::Wrapper - # @returns [Boolean] Always false, so that reads continue until the stream ends or fails. - def empty? - false - end - end - - private_constant :ResponseBody - # Build a compatible stub class for a generated GRPC::GenericService. # @parameter service [Class] The generated service definition. # @returns [Class] A client stub with methods for the service's unary RPCs. @@ -327,7 +317,6 @@ def perform_request(request, operation) begin operation.metadata = extract_metadata(Protocol::HTTP::Headers.new(response.headers.header.to_a, policy: Protocol::GRPC::HEADER_POLICY)) response_encoding = response.headers["grpc-encoding"] - response.body = ResponseBody.new(response.body) if response.body response_body = Protocol::GRPC::Body::Readable.wrap(response, encoding: response_encoding) payload = response_body&.read response_body&.finish