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
2 changes: 1 addition & 1 deletion async-grpc-compatible.gemspec
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,6 @@ Gem::Specification.new do |specification|
specification.add_dependency "async-grpc", "~> 0.10"
specification.add_dependency "async-http", "~> 0.100"
specification.add_dependency "grpc"
specification.add_dependency "io-endpoint", "~> 0.18"
specification.add_dependency "io-endpoint", "~> 0.19"
specification.add_dependency "protocol-grpc", "~> 0.17"
end
98 changes: 84 additions & 14 deletions lib/async/grpc/compatible/client_stub.rb
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,8 @@
require_relative "channel_credentials"
require_relative "operation"

::Thread.attr_accessor :async_grpc_compatible_shared_clients

module Async
module GRPC
module Compatible
Expand Down Expand Up @@ -45,9 +47,50 @@ def close
end
end

# Represents a channel whose client is shared with other shared channels for the same URL and TLS configuration on the same thread.
#
# Each thread has its own client because an Async connection pool belongs to a single reactor.
class SharedChannel < Channel
# @returns [Hash] The current thread's shared clients, keyed by URL and TLS configuration.
def self.clients
::Thread.current.async_grpc_compatible_shared_clients ||= {}
end

# Close the current thread's shared clients. Shared channels open new clients on their next call.
def self.close
# Detach the clients before closing them, since closing may yield to a call which opens a new client:
clients = ::Thread.current.async_grpc_compatible_shared_clients
::Thread.current.async_grpc_compatible_shared_clients = nil

clients&.each_value(&:close)
end

# Initialize a shared channel for the given URL and TLS configuration.
#
# The channel freezes a copy of the TLS configuration, so later changes by the caller do not affect shared clients.
#
# @parameter url [String] The remote `http` or `https` URL.
# @parameter tls_configuration [IO::Endpoint::TLS::Configuration | Nil] The TLS configuration for an `https` URL.
def initialize(url, tls_configuration = nil)
tls_configuration = tls_configuration&.dup&.freeze
@endpoint = Async::HTTP::Endpoint.parse(url, protocol: Async::HTTP::Protocol::HTTP2, tls_configuration: tls_configuration)
@key = ["#{@endpoint.scheme}://#{@endpoint.authority}", tls_configuration]
end

# @attribute [Async::GRPC::Client] The current thread's client for this endpoint.
def client
self.class.clients[@key] ||= Async::GRPC::Client.open(@endpoint)
end

# Leave the shared client open for other channels.
def close
end
end

# Represents a subset of `GRPC::ClientStub` backed by {Async::GRPC::Client}.
class ClientStub
INSECURE_CREDENTIALS = :this_channel_is_insecure
LOCAL_SUBCHANNEL_POOL = "grpc.use_local_subchannel_pool"
DEFAULT_TIMEOUT = nil

# Transport failures which end a call without a gRPC status. grpc-ruby reports these as `UNAVAILABLE`.
Expand Down Expand Up @@ -88,6 +131,9 @@ def self.for(service)
end

# Construct a compatible channel.
#
# Without an override, the channel shares connections with other stubs for the same target and TLS configuration on the current thread. Set the `grpc.use_local_subchannel_pool` channel argument to a non-zero value to use a client owned by this channel instead.
#
# @parameter channel_override [Channel, Async::GRPC::Client | Nil] An existing compatible channel or client.
# @parameter host [String] The gRPC target.
# @parameter credentials [IO::Endpoint::TLS::Configuration, Symbol] The channel TLS configuration or insecure marker.
Expand All @@ -105,16 +151,35 @@ def self.setup_channel(channel_override, host, credentials, channel_arguments =
raise TypeError, "Channel override must be an Async::GRPC::Compatible::Channel or Async::GRPC::Client!"
end

endpoint = endpoint_for(host, credentials, channel_arguments)
Channel.new(endpoint)
# grpc-ruby accepts string and symbol keys. Only fall back on nil, so an explicit false is preserved:
local_pool = channel_arguments[LOCAL_SUBCHANNEL_POOL]
local_pool = channel_arguments[LOCAL_SUBCHANNEL_POOL.to_sym] if local_pool.nil?

# gRPC uses integer boolean flags (0/1). Ruby treats 0 as truthy, so check it explicitly:
if local_pool && local_pool != 0
Channel.new(endpoint_for(host, credentials, channel_arguments))
else
SharedChannel.new(url_for(host, credentials), tls_configuration_for(credentials))
end
end

# Construct an HTTP/2 endpoint for a gRPC target.
#
# Channels with a local subchannel pool use this endpoint. Shared channels are identified by their URL and TLS configuration, so they are constructed from {url_for} and {tls_configuration_for} instead.
#
# @parameter host [String] The gRPC target.
# @parameter credentials [IO::Endpoint::TLS::Configuration, Symbol] The channel TLS configuration or insecure marker.
# @parameter channel_arguments [Hash] gRPC channel arguments.
# @returns [Async::HTTP::Endpoint] The HTTP/2 endpoint.
def self.endpoint_for(host, credentials, channel_arguments = {})
Async::HTTP::Endpoint.parse(url_for(host, credentials), protocol: Async::HTTP::Protocol::HTTP2, tls_configuration: tls_configuration_for(credentials))
end

# Construct the URL for a gRPC target.
# @parameter host [String] The gRPC target.
# @parameter credentials [IO::Endpoint::TLS::Configuration, Symbol] The channel TLS configuration or insecure marker.
# @returns [String] The `http` or `https` URL.
def self.url_for(host, credentials)
raise TypeError, "Host must be a String!" unless host.is_a?(String)

scheme = scheme_for(credentials)
Expand All @@ -126,18 +191,23 @@ def self.endpoint_for(host, credentials, channel_arguments = {})
url = "#{scheme}://#{target}"
end

if scheme == "https"
configuration = IO::Endpoint::TLS::Configuration.new(
trust_store: credentials.trust_store,
certificate_chain: credentials.certificate_chain,
private_key: credentials.private_key,
verification: credentials.verification || :peer
)
end
raise ArgumentError, "Target scheme must match the channel credentials!" unless url.start_with?("#{scheme}://")

endpoint = Async::HTTP::Endpoint.parse(url, protocol: Async::HTTP::Protocol::HTTP2, tls_configuration: configuration)
raise ArgumentError, "Target scheme must match the channel credentials!" unless endpoint.scheme == scheme
endpoint
return url
end

# Construct the TLS configuration for the given credentials.
# @parameter credentials [IO::Endpoint::TLS::Configuration, Symbol] The channel TLS configuration or insecure marker.
# @returns [IO::Endpoint::TLS::Configuration | Nil] The TLS configuration, or `nil` for insecure credentials.
def self.tls_configuration_for(credentials)
return nil if scheme_for(credentials) == "http"

IO::Endpoint::TLS::Configuration.new(
trust_store: credentials.trust_store,
certificate_chain: credentials.certificate_chain,
private_key: credentials.private_key,
verification: credentials.verification || :peer
)
end

# Determine the URL scheme for the given credentials.
Expand Down Expand Up @@ -234,7 +304,7 @@ def request_response(method, request, marshal, unmarshal,
operation.execute
end

# Close a channel created by this stub.
# Close a channel created by this stub. Shared clients remain open for other stubs.
def close
@channel.close if @owned_channel
end
Expand Down
2 changes: 1 addition & 1 deletion lib/async/grpc/compatible/gapic.rb
Original file line number Diff line number Diff line change
Expand Up @@ -39,7 +39,7 @@ def create_channel_pool(...)
raise ArgumentError, "Use a shared Async channel instead of a GAPIC channel pool!"
end

# Close the underlying stub's owned connection pool.
# Close the underlying stub's owned connection pool, leaving shared clients open.
def close
@grpc_stub&.close
end
Expand Down
24 changes: 21 additions & 3 deletions readme.md
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@

grpc-ruby compatible client interfaces backed by `async-grpc` and `async-http`.

The gem is intended for generated clients which currently construct a `GRPC::ClientStub`, but need to make non-blocking calls inside an Async event loop. Connection reuse and HTTP/2 multiplexing remain the responsibility of `async-http`; this gem does not add a second connection pool.
The gem is intended for generated clients which currently construct a `GRPC::ClientStub`, but need to make non-blocking calls inside an Async event loop. Connection pooling and HTTP/2 multiplexing remain the responsibility of `async-http`; stubs for the same target share its connection pool, as described in [Connection sharing](#connection-sharing).

The gem depends on `grpc` for service definitions and error types. TLS configuration comes from `IO::Endpoint`, and requests use Async's connection pool.

Expand Down Expand Up @@ -48,18 +48,36 @@ The initial implementation supports:
- Deferred unary operations using `return_op: true`.
- Ruby credential updaters supplied through `call_credentials:`, `credentials:`, or a credential object with `updater_proc`.
- Unary stub generation from `GRPC::GenericService` definitions.
- Connection sharing between stubs for the same target, with `grpc.use_local_subchannel_pool` to opt out.

The following are not yet supported:

- Client, server, or bidirectional streaming.
- grpc-ruby interceptors.
- Parent call propagation.
- Native `GRPC::Core::ChannelCredentials`, `GRPC::Core::CallCredentials`, composed credentials, and native channel overrides. These are rejected because their TLS configuration and authentication callbacks cannot be recovered through Ruby's public API.
- grpc-ruby channel arguments beyond accepting the compatible constructor parameter.
- grpc-ruby channel arguments other than `grpc.use_local_subchannel_pool`.
- Non-DNS resolvers such as Unix sockets and xDS.

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.

## Connection sharing

Stubs constructed without a channel override share connections, similar to grpc-core's global subchannel pool. Stubs on the same thread use one `Async::GRPC::Client` when they connect to the same target with the same TLS configuration. TLS configurations are compared by value, so separately constructed but identical configurations share a client, and different trust roots, client certificates, or verification policies never do. Constructing a stub or GAPIC client per request or per job therefore reuses warm connections within the same Async reactor instead of performing a new TLS and HTTP/2 handshake.

Each thread has its own shared client, because an Async connection pool belongs to a single reactor. The client is selected when each call runs, so a stub shared between threads uses the calling thread's connections. Connections do not outlive their reactor. Calls made outside a reactor each run in a temporary reactor, so run related calls inside one `Sync` or `Async` block to reuse connections.

Closing a stub leaves shared clients open for other stubs. Call `Async::GRPC::Compatible::SharedChannel.close` to close the current thread's shared clients, for example when a worker thread shuts down; later calls open new clients. To give a stub a connection pool of its own, which `close` releases, set grpc's local subchannel pool argument:

``` ruby
stub = Async::GRPC::Compatible::ClientStub.new(
"grpc.example.com:443", IO::Endpoint::TLS::Configuration.new,
channel_args: {"grpc.use_local_subchannel_pool" => 1}
)
```

Pass `channel_override:` (or `channel:` to `GapicServiceStub`) to control connection reuse explicitly. The caller owns a channel supplied this way. Shared channels are built from the target URL and TLS configuration alone, so subclasses which override `ClientStub.endpoint_for`, for example to connect through a custom transport, only take effect with a local subchannel pool; alternatively, pass `channel_override: Async::GRPC::Compatible::Channel.new(endpoint)`.

## 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 Expand Up @@ -154,7 +172,7 @@ Sync do
end
```

For an existing Async connection pool, pass `channel:` to the adapter. The caller owns that channel. GAPIC native channel pools are unsupported because Async::HTTP already manages connections.
Without `channel:`, adapters share connections as described in [Connection sharing](#connection-sharing). For an existing Async connection pool, pass `channel:` to the adapter. The caller owns that channel. GAPIC native channel pools are unsupported because Async::HTTP already manages connections.

Generated high-level Google clients construct `Gapic::ServiceStub` inside their constructors. Applications adapting those constructors can use `GapicServiceStub.new(Service, credentials: original_credentials, ...)` at that construction point. The adapter can also be called directly, as above, without replacing global GRPC constants. Keep the original Ruby credentials available at this boundary; credentials already composed into `GRPC::Core::ChannelCredentials` cannot be recovered.

Expand Down
1 change: 1 addition & 0 deletions releases.md
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
- 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.
- Share connections between stubs for the same target and TLS configuration on each thread, like grpc-core's global subchannel pool. Closing a stub leaves shared clients open; use `SharedChannel.close` to close the current thread's shared clients, or set `grpc.use_local_subchannel_pool` for a stub-owned connection pool.

## v0.0.0

Expand Down
Loading
Loading