From cbde3d094d92f26740793371f1248aab43a77fb7 Mon Sep 17 00:00:00 2001 From: Martin Ek Date: Thu, 24 Sep 2026 12:59:21 -0700 Subject: [PATCH 1/4] Share connections between stubs for the same target. Stubs constructed without a channel override now share one client per thread for the same URL and TLS configuration, similar to grpc-core's global subchannel pool. TLS configurations are copied and compared by value, and cache keys contain a digest rather than private keys. Closing a stub leaves shared clients open; `SharedChannel.close` closes the current thread's clients, and `grpc.use_local_subchannel_pool` gives a stub its own client. Fixes #8. Co-Authored-By: Claude Opus 5.5 (1M context) --- lib/async/grpc/compatible/client_stub.rb | 96 ++++++++++-- lib/async/grpc/compatible/gapic.rb | 2 +- readme.md | 24 ++- releases.md | 1 + test/async/grpc/compatible/client_stub.rb | 183 ++++++++++++++++++++++ test/async/grpc/compatible/gapic.rb | 13 +- 6 files changed, 298 insertions(+), 21 deletions(-) diff --git a/lib/async/grpc/compatible/client_stub.rb b/lib/async/grpc/compatible/client_stub.rb index dfefadc..2342978 100644 --- a/lib/async/grpc/compatible/client_stub.rb +++ b/lib/async/grpc/compatible/client_stub.rb @@ -10,12 +10,15 @@ require "openssl" require "grpc" require "io/endpoint/tls/configuration" +require "openssl" require "protocol/grpc/body/readable" require "protocol/grpc/body/writable" require "protocol/grpc/metadata" require_relative "channel_credentials" require_relative "operation" +::Thread.attr_accessor :async_grpc_compatible_shared_clients + module Async module GRPC module Compatible @@ -45,9 +48,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 digest. + 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 copies the TLS configuration, so later changes by the caller do not affect shared clients. Its key contains a digest of the TLS configuration rather than its private key. + # + # @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) + data = Marshal.dump(tls_configuration) + @endpoint = Async::HTTP::Endpoint.parse(url, protocol: Async::HTTP::Protocol::HTTP2, tls_configuration: Marshal.load(data)) + @key = ["#{@endpoint.scheme}://#{@endpoint.authority}", OpenSSL::Digest::SHA256.hexdigest(data)] + 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`. @@ -88,6 +132,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. @@ -105,16 +152,33 @@ 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) + local_pool = channel_arguments[LOCAL_SUBCHANNEL_POOL] + local_pool = channel_arguments[LOCAL_SUBCHANNEL_POOL.to_sym] if local_pool.nil? + + 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) @@ -126,18 +190,22 @@ 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}://") + 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" - 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 + 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. @@ -234,7 +302,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 diff --git a/lib/async/grpc/compatible/gapic.rb b/lib/async/grpc/compatible/gapic.rb index 0351a8c..07cc588 100644 --- a/lib/async/grpc/compatible/gapic.rb +++ b/lib/async/grpc/compatible/gapic.rb @@ -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 diff --git a/readme.md b/readme.md index 53d236f..e73bcff 100644 --- a/readme.md +++ b/readme.md @@ -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. @@ -48,6 +48,7 @@ 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: @@ -55,11 +56,28 @@ The following are not yet supported: - 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. @@ -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. diff --git a/releases.md b/releases.md index 1600352..adf442f 100644 --- a/releases.md +++ b/releases.md @@ -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 diff --git a/test/async/grpc/compatible/client_stub.rb b/test/async/grpc/compatible/client_stub.rb index 311d970..3d4e831 100644 --- a/test/async/grpc/compatible/client_stub.rb +++ b/test/async/grpc/compatible/client_stub.rb @@ -293,6 +293,189 @@ def request(value, **options) direct_stub&.close end + with "shared channels" do + def echo(stub, value) + stub.request_response("/#{service_name}/Echo", CompatibleMessage.new(value), CompatibleMessage.method(:encode), CompatibleMessage.method(:decode)).value + end + + it "applies an endpoint override with a local subchannel pool" do + port = URI(bound_url).port + mapped_class = Class.new(subject) do + define_singleton_method(:endpoint_for) do |host, credentials, channel_arguments = {}| + endpoint = super(host, credentials, channel_arguments) + endpoint.endpoint = IO::Endpoint.tcp("127.0.0.1", port) + endpoint + end + end + mapped = mapped_class.new("127.0.0.1:1", :this_channel_is_insecure, channel_args: {"grpc.use_local_subchannel_pool" => 1}) + + expect(echo(mapped, "mapped")).to be == "mapped" + ensure + mapped&.close + end + + it "validates the target scheme" do + expect{subject.new("http://example.com", IO::Endpoint::TLS::Configuration.new)}.to raise_exception(ArgumentError, message: be =~ /scheme/) + end + + it "is unaffected by later changes to its URL" do + url = bound_url.dup + channel = Async::GRPC::Compatible::SharedChannel.new(url) + client = channel.client + + url.replace("http://other.invalid:1") + + expect(channel.client).to be_equal(client) + expect(subject.new(bound_url, :this_channel_is_insecure).channel.client).to be_equal(client) + expect(echo(subject.new("unused", nil, channel_override: channel), "unchanged")).to be == "unchanged" + end + + it "shares one connection between stubs for the same target" do + first = subject.new(bound_url, :this_channel_is_insecure) + second = subject.new(bound_url, :this_channel_is_insecure) + + expect(first.channel).to be_a(Async::GRPC::Compatible::SharedChannel) + expect(second.channel.client).to be_equal(first.channel.client) + expect(echo(first, "first")).to be == "first" + expect(echo(second, "second")).to be == "second" + expect(first.channel.client.delegate.pool.size).to be == 1 + end + + it "keeps the shared client open when a stub closes" do + first = subject.new(bound_url, :this_channel_is_insecure) + expect(echo(first, "first")).to be == "first" + first.close + + second = subject.new(bound_url, :this_channel_is_insecure) + expect(echo(second, "second")).to be == "second" + expect(second.channel.client.delegate.pool.size).to be == 1 + end + + it "closes the current thread's shared clients" do + shared = subject.new(bound_url, :this_channel_is_insecure) + expect(echo(shared, "before")).to be == "before" + client = shared.channel.client + expect(client.delegate.pool.size).to be == 1 + + Async::GRPC::Compatible::SharedChannel.close + + expect(client.delegate.pool.size).to be == 0 + expect(shared.channel.client).not.to be_equal(client) + expect(echo(shared, "after")).to be == "after" + end + + it "uses a separate client on each thread" do + shared = subject.new(bound_url, :this_channel_is_insecure) + other = Thread.new{shared.channel.client}.value + + expect(other).not.to be_equal(shared.channel.client) + expect(Thread.new{shared.channel.client}.value).not.to be_equal(other) + end + + it "uses an owned client for a local subchannel pool" do + shared = subject.new(bound_url, :this_channel_is_insecure) + local = subject.new(bound_url, :this_channel_is_insecure, channel_args: {"grpc.use_local_subchannel_pool" => 1}) + + expect(local.channel).not.to be_a(Async::GRPC::Compatible::SharedChannel) + expect(local.channel.client).not.to be_equal(shared.channel.client) + expect(echo(local, "local")).to be == "local" + + expect(local.channel.client).to receive(:close) + local.close + end + + it "accepts a symbol key for the local subchannel pool" do + local = subject.new(bound_url, :this_channel_is_insecure, channel_args: {"grpc.use_local_subchannel_pool": 1}) + + expect(local.channel).not.to be_a(Async::GRPC::Compatible::SharedChannel) + end + + it "shares the client when the local subchannel pool is disabled" do + shared = subject.new(bound_url, :this_channel_is_insecure) + disabled = subject.new(bound_url, :this_channel_is_insecure, channel_args: {"grpc.use_local_subchannel_pool" => 0}) + + expect(disabled.channel.client).to be_equal(shared.channel.client) + end + + with "TLS" do + include TLSContext + + it "shares a client between equivalent TLS configurations" do + first = subject.new(bound_url, tls_credentials) + second = subject.new(bound_url, tls_credentials) + + expect(second.channel.client).to be_equal(first.channel.client) + expect(echo(first, "first")).to be == "first" + expect(echo(second, "second")).to be == "second" + expect(first.channel.client.delegate.pool.size).to be == 1 + end + + it "does not share a client between different TLS configurations" do + trusted = subject.new(bound_url, tls_credentials) + untrusted = subject.new(bound_url, IO::Endpoint::TLS::Configuration.new) + + expect(untrusted.channel.client).not.to be_equal(trusted.channel.client) + expect(echo(trusted, "trusted")).to be == "trusted" + expect{echo(untrusted, "untrusted")}.to raise_exception(::GRPC::Unavailable, message: be =~ /certificate verify failed/).and(have_attributes( + cause: be_a(OpenSSL::SSL::SSLError) + )) + end + + it "does not reuse a client after the caller changes its TLS configuration" do + credentials = tls_credentials + original = subject.new(bound_url, credentials) + expect(echo(original, "original")).to be == "original" + + replacement_authority = Object.new.extend(Sus::Fixtures::OpenSSL::CertificateAuthorityContext) + credentials.trust_store.certificates.replace([replacement_authority.certificate_authority_certificate.to_pem]) + + expect(original.channel.endpoint.tls_configuration.trust_store.certificates).to be == tls_credentials.trust_store.certificates + expect(subject.new(bound_url, tls_credentials).channel.client).to be_equal(original.channel.client) + + replacement = subject.new(bound_url, credentials) + expect(replacement.channel.client).not.to be_equal(original.channel.client) + expect{echo(replacement, "replacement")}.to raise_exception(::GRPC::Unavailable, message: be =~ /certificate verify failed/).and(have_attributes( + cause: be_a(OpenSSL::SSL::SSLError) + )) + end + + it "does not expose private keys when inspected" do + credentials = IO::Endpoint::TLS::Configuration.new( + trust_store: tls_credentials.trust_store, + certificate_chain: [certificate.to_pem], private_key: key.to_pem + ) + identified = subject.new(bound_url, credentials) + identified.channel.client + private_key = key.to_pem.lines[1].chomp + + expect(identified.inspect).not.to be(:include?, private_key) + expect(Async::GRPC::Compatible::SharedChannel.clients.inspect).not.to be(:include?, private_key) + end + + it "does not share a client between different client certificates" do + anonymous = subject.new(bound_url, tls_credentials) + identified = subject.new(bound_url, IO::Endpoint::TLS::Configuration.new( + trust_store: tls_credentials.trust_store, + certificate_chain: [certificate.to_pem], private_key: key.to_pem + )) + + expect(identified.channel.client).not.to be_equal(anonymous.channel.client) + end + + it "shares a client between GAPIC stubs" do + credentials = ->{Async::GRPC::Compatible::ChannelCredentials.new(certificate_authority_certificate.to_pem)} + first = Async::GRPC::Compatible::GapicServiceStub.new(GeneratedCompatibleService, endpoint: bound_url, credentials: credentials.call, logger: nil) + second = Async::GRPC::Compatible::GapicServiceStub.new(GeneratedCompatibleService, endpoint: bound_url, credentials: credentials.call, logger: nil) + + expect(first.call_rpc(:echo, CompatibleMessage.new("first")).value).to be == "first" + first.close + expect(second.call_rpc(:echo, CompatibleMessage.new("second")).value).to be == "second" + expect(second.grpc_stub.channel.client).to be_equal(first.grpc_stub.channel.client) + expect(second.grpc_stub.channel.client.delegate.pool.size).to be == 1 + end + end + end + it "does not block sibling fibers" do client_stub = stub events = [] diff --git a/test/async/grpc/compatible/gapic.rb b/test/async/grpc/compatible/gapic.rb index c43f566..578cfbe 100644 --- a/test/async/grpc/compatible/gapic.rb +++ b/test/async/grpc/compatible/gapic.rb @@ -32,9 +32,16 @@ stub.close end - it "closes the client when it owns the channel" do - owned_stub = subject.new(service, endpoint: "example.googleapis.com", credentials: credentials, logger: nil) - expect(owned_stub.grpc_stub.channel.endpoint.to_url.to_s).to be == "https://example.googleapis.com/" + it "leaves the default shared client open when closed" do + shared_stub = subject.new(service, endpoint: "example.googleapis.com", credentials: credentials, logger: nil) + expect(shared_stub.grpc_stub.channel.endpoint.to_url.to_s).to be == "https://example.googleapis.com/" + expect(shared_stub.grpc_stub.channel.client).not.to receive(:close) + + shared_stub.close + end + + it "closes the client when it uses a local subchannel pool" do + owned_stub = subject.new(service, endpoint: "example.googleapis.com", credentials: credentials, channel_args: {"grpc.use_local_subchannel_pool" => 1}, logger: nil) expect(owned_stub.grpc_stub.channel.client).to receive(:close) owned_stub.close From 068939e0eeca47c091a994952c2b3947adc6f8f3 Mon Sep 17 00:00:00 2001 From: Samuel Williams Date: Fri, 25 Sep 2026 10:54:02 +1200 Subject: [PATCH 2/4] Use TLS configuration value equality for shared clients --- async-grpc-compatible.gemspec | 2 +- lib/async/grpc/compatible/client_stub.rb | 11 +++++------ 2 files changed, 6 insertions(+), 7 deletions(-) diff --git a/async-grpc-compatible.gemspec b/async-grpc-compatible.gemspec index 93bc714..418805a 100644 --- a/async-grpc-compatible.gemspec +++ b/async-grpc-compatible.gemspec @@ -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 diff --git a/lib/async/grpc/compatible/client_stub.rb b/lib/async/grpc/compatible/client_stub.rb index 2342978..77343f3 100644 --- a/lib/async/grpc/compatible/client_stub.rb +++ b/lib/async/grpc/compatible/client_stub.rb @@ -10,7 +10,6 @@ require "openssl" require "grpc" require "io/endpoint/tls/configuration" -require "openssl" require "protocol/grpc/body/readable" require "protocol/grpc/body/writable" require "protocol/grpc/metadata" @@ -52,7 +51,7 @@ def close # # 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 digest. + # @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 @@ -68,14 +67,14 @@ def self.close # Initialize a shared channel for the given URL and TLS configuration. # - # The channel copies the TLS configuration, so later changes by the caller do not affect shared clients. Its key contains a digest of the TLS configuration rather than its private key. + # 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) - data = Marshal.dump(tls_configuration) - @endpoint = Async::HTTP::Endpoint.parse(url, protocol: Async::HTTP::Protocol::HTTP2, tls_configuration: Marshal.load(data)) - @key = ["#{@endpoint.scheme}://#{@endpoint.authority}", OpenSSL::Digest::SHA256.hexdigest(data)] + 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. From 3ab755d8eccb541b0aa7a80bfcfe89f7fc589ab7 Mon Sep 17 00:00:00 2001 From: Samuel Williams Date: Fri, 25 Sep 2026 10:58:14 +1200 Subject: [PATCH 3/4] Return the constructed URL explicitly --- lib/async/grpc/compatible/client_stub.rb | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/lib/async/grpc/compatible/client_stub.rb b/lib/async/grpc/compatible/client_stub.rb index 77343f3..74e89fa 100644 --- a/lib/async/grpc/compatible/client_stub.rb +++ b/lib/async/grpc/compatible/client_stub.rb @@ -190,7 +190,8 @@ def self.url_for(host, credentials) end raise ArgumentError, "Target scheme must match the channel credentials!" unless url.start_with?("#{scheme}://") - url + + return url end # Construct the TLS configuration for the given credentials. From 77b9322114907844191c52e2c016e7d394c59e3c Mon Sep 17 00:00:00 2001 From: Samuel Williams Date: Fri, 25 Sep 2026 11:11:08 +1200 Subject: [PATCH 4/4] Explain gRPC channel argument compatibility --- lib/async/grpc/compatible/client_stub.rb | 2 ++ 1 file changed, 2 insertions(+) diff --git a/lib/async/grpc/compatible/client_stub.rb b/lib/async/grpc/compatible/client_stub.rb index 74e89fa..7022938 100644 --- a/lib/async/grpc/compatible/client_stub.rb +++ b/lib/async/grpc/compatible/client_stub.rb @@ -151,9 +151,11 @@ 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 + # 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