Skip to content
Merged
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
1 change: 1 addition & 0 deletions async-grpc-compatible.gemspec
Original file line number Diff line number Diff line change
Expand Up @@ -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
37 changes: 36 additions & 1 deletion lib/async/grpc/compatible/client_stub.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -49,6 +50,25 @@ 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

# 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.
Expand Down Expand Up @@ -285,6 +305,13 @@ 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
Expand All @@ -302,10 +329,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)
Expand Down
2 changes: 0 additions & 2 deletions readme.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
1 change: 1 addition & 0 deletions releases.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
146 changes: 143 additions & 3 deletions test/async/grpc/compatible/client_stub.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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"
Expand Down Expand Up @@ -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
Expand All @@ -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
Expand Down Expand Up @@ -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
Expand Down
Loading