diff --git a/async-grpc-compatible.gemspec b/async-grpc-compatible.gemspec index c81d9e0..1f33bd9 100644 --- a/async-grpc-compatible.gemspec +++ b/async-grpc-compatible.gemspec @@ -21,6 +21,8 @@ Gem::Specification.new do |specification| specification.required_ruby_version = ">= 3.3" - specification.add_dependency "async-grpc", "~> 0.8" + 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" end diff --git a/gems.rb b/gems.rb index 6f32141..299297b 100644 --- a/gems.rb +++ b/gems.rb @@ -26,6 +26,9 @@ end group :test do + gem "gapic-common" + gem "googleauth" + gem "covered" gem "sus" @@ -34,6 +37,7 @@ gem "rubocop-socketry" gem "sus-fixtures-async-http" + gem "sus-fixtures-openssl" gem "bake-test" gem "bake-test-external" diff --git a/lib/async/grpc/compatible/channel_credentials.rb b/lib/async/grpc/compatible/channel_credentials.rb new file mode 100644 index 0000000..99daec2 --- /dev/null +++ b/lib/async/grpc/compatible/channel_credentials.rb @@ -0,0 +1,39 @@ +# frozen_string_literal: true + +# Released under the MIT License. +# Copyright, 2026, by Samuel Williams. + +require "io/endpoint/tls/configuration" + +module Async + module GRPC + module Compatible + # Maps gRPC channel credentials to transport-neutral TLS configurations. + module ChannelCredentials + # Create a TLS configuration using the positional arguments of `GRPC::Core::ChannelCredentials.new`. + # @parameter root_certificates [String | Nil] The trusted root certificates encoded as a PEM bundle, or nil to use the transport's default trust store. + # @parameter private_key [String | Nil] The client private key encoded as PEM. + # @parameter certificate_chain [String | Nil] The client certificate chain encoded as a PEM bundle, with the leaf certificate first. + # @returns [IO::Endpoint::TLS::Configuration] The transport-neutral TLS configuration. + # @raises [ArgumentError] If a certificate bundle is empty or the client certificate chain and private key are not supplied together. + # @raises [TypeError] If certificate or private key material is not a string. + def self.new(root_certificates = nil, private_key = nil, certificate_chain = nil) + trust_store = unless root_certificates.nil? + IO::Endpoint::TLS::TrustStore.parse(root_certificates) + end + + certificates = unless certificate_chain.nil? + IO::Endpoint::TLS::Certificates.parse(certificate_chain) + end + + IO::Endpoint::TLS::Configuration.new( + trust_store: trust_store, + certificate_chain: certificates, + private_key: private_key, + verification: :peer + ) + end + end + end + end +end diff --git a/lib/async/grpc/compatible/client_stub.rb b/lib/async/grpc/compatible/client_stub.rb index 5202861..8862bc7 100644 --- a/lib/async/grpc/compatible/client_stub.rb +++ b/lib/async/grpc/compatible/client_stub.rb @@ -8,9 +8,12 @@ require "async/http/protocol/http2" require "base64" require "grpc" +require "io/endpoint/tls/configuration" require "protocol/grpc/body/readable" require "protocol/grpc/body/writable" require "protocol/grpc/metadata" +require_relative "channel_credentials" +require_relative "operation" module Async module GRPC @@ -18,10 +21,13 @@ module Compatible # Represents a reusable Async gRPC channel. class Channel # Initialize a channel for the given endpoint. - # @parameter endpoint [Async::HTTP::Endpoint] The remote HTTP/2 endpoint. + # @parameter endpoint [Async::HTTP::Endpoint | Nil] The remote HTTP/2 endpoint, inferred from a supplied Async client when possible. # @parameter client [Async::GRPC::Client | Nil] An existing client to use. def initialize(endpoint = nil, client: nil) @endpoint = endpoint + if @endpoint.nil? && client.is_a?(Async::GRPC::Client) && client.delegate.respond_to?(:endpoint) + @endpoint = client.delegate.endpoint + end @client = client || Async::GRPC::Client.open(endpoint) @owned = client.nil? end @@ -43,10 +49,28 @@ class ClientStub INSECURE_CREDENTIALS = :this_channel_is_insecure DEFAULT_TIMEOUT = nil + # 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. + def self.for(service) + Class.new(self) do + service.rpc_descs.each do |name, description| + method_name = ::GRPC::GenericService.underscore(name.to_s) + path = "/#{service.service_name}/#{name}" + marshal = description.marshal_proc + unmarshal = description.unmarshal_proc(:output) + define_method(method_name) do |request, **options| + raise NotImplementedError, "Streaming RPCs are not yet supported!" unless description.request_response? + request_response(path, request, marshal, unmarshal, **options) + end + end + end + end + # Construct a compatible channel. # @parameter channel_override [Channel, Async::GRPC::Client | Nil] An existing compatible channel or client. # @parameter host [String] The gRPC target. - # @parameter credentials [GRPC::Core::ChannelCredentials, Symbol] The channel credentials. + # @parameter credentials [IO::Endpoint::TLS::Configuration, Symbol] The channel TLS configuration or insecure marker. # @parameter channel_arguments [Hash] gRPC channel arguments. # @returns [Channel] The compatible channel. def self.setup_channel(channel_override, host, credentials, channel_arguments = {}) @@ -58,7 +82,7 @@ def self.setup_channel(channel_override, host, credentials, channel_arguments = when nil # Continue constructing the channel: else - raise TypeError, "channel_override must be an Async::GRPC::Compatible::Channel or Async::GRPC::Client" + raise TypeError, "Channel override must be an Async::GRPC::Compatible::Channel or Async::GRPC::Client!" end endpoint = endpoint_for(host, credentials, channel_arguments) @@ -67,11 +91,11 @@ def self.setup_channel(channel_override, host, credentials, channel_arguments = # Construct an HTTP/2 endpoint for a gRPC target. # @parameter host [String] The gRPC target. - # @parameter credentials [GRPC::Core::ChannelCredentials, Symbol] The channel credentials. + # @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 = {}) - raise TypeError, "host must be a String" unless host.is_a?(String) + raise TypeError, "Host must be a String!" unless host.is_a?(String) scheme = scheme_for(credentials) target = normalize_target(host) @@ -82,20 +106,31 @@ def self.endpoint_for(host, credentials, channel_arguments = {}) url = "#{scheme}://#{target}" end - Async::HTTP::Endpoint.parse(url, protocol: Async::HTTP::Protocol::HTTP2) + 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 + + 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 end # Determine the URL scheme for the given credentials. - # @parameter credentials [GRPC::Core::ChannelCredentials, Symbol] The channel credentials. + # @parameter credentials [IO::Endpoint::TLS::Configuration, Symbol] The channel TLS configuration or insecure marker. # @returns [String] Either `"http"` or `"https"`. def self.scheme_for(credentials) return "http" if credentials == INSECURE_CREDENTIALS - if credentials.is_a?(::GRPC::Core::ChannelCredentials) + if credentials.is_a?(IO::Endpoint::TLS::Configuration) return "https" end - raise TypeError, "credentials must be GRPC channel credentials or :this_channel_is_insecure" + raise TypeError, "Credentials must be IO::Endpoint::TLS::Configuration or :this_channel_is_insecure; native gRPC credentials are unsupported!" end # Normalize a grpc-ruby target into an HTTP authority. @@ -107,7 +142,7 @@ def self.normalize_target(host) elsif host.start_with?("dns://") host.delete_prefix("dns://").delete_prefix("/") elsif host.match?(/\A(?:unix|unix-abstract|ipv4|ipv6|xds|passthrough):/i) - raise ArgumentError, "Unsupported gRPC target: #{host.inspect}" + raise ArgumentError, "Unsupported gRPC target: #{host.inspect}!" else host end @@ -115,20 +150,29 @@ def self.normalize_target(host) # Create a compatible client stub. # @parameter host [String] The gRPC target. - # @parameter credentials [GRPC::Core::ChannelCredentials, Symbol, Nil] The channel credentials. + # @parameter credentials [IO::Endpoint::TLS::Configuration, Symbol, Proc, Object, Nil] The channel TLS configuration, insecure marker, or Ruby authentication callback. Callbacks use a default verified TLS channel. Nil requires a channel override. # @parameter channel_override [Channel, Async::GRPC::Client | Nil] An existing compatible channel or client. # @parameter timeout [Numeric | Nil] The default relative timeout in seconds. # @parameter propagate_mask [Integer | Nil] Reserved for grpc-ruby compatibility. # @parameter channel_args [Hash] gRPC channel arguments. + # @parameter call_credentials [Proc | Object | Nil] An authentication callback or an object with updater_proc. # @parameter interceptors [Array] grpc-ruby client interceptors, which are not yet supported. def initialize(host, credentials, channel_override: nil, timeout: nil, propagate_mask: nil, channel_args: {}, - interceptors: []) - raise NotImplementedError, "Client interceptors are not yet supported" unless interceptors.empty? + interceptors: [], + call_credentials: nil) + raise NotImplementedError, "Client interceptors are not yet supported!" unless interceptors.empty? + if credentials.respond_to?(:updater_proc) || credentials.respond_to?(:call) + raise ArgumentError, "Supply call credentials only once!" if call_credentials + call_credentials = credentials + credentials = IO::Endpoint::TLS::Configuration.new(verification: :peer) + end + self.class.scheme_for(credentials) unless credentials.nil? && channel_override + @call_credentials = call_credentials channel_arguments = channel_args.dup @channel = self.class.setup_channel(channel_override, host, credentials, channel_arguments) @owned_channel = channel_override.nil? @@ -158,57 +202,104 @@ def request_response(method, request, marshal, unmarshal, parent: nil, credentials: nil, metadata: {}) - raise NotImplementedError, "return_op is not yet supported" if return_op - raise NotImplementedError, "parent call propagation is not yet supported" if parent - raise NotImplementedError, "per-call credentials are not yet supported" if credentials + raise NotImplementedError, "Parent call propagation is not yet supported!" if parent timeout = relative_timeout(deadline) + call_deadline = timeout && Time.now + timeout + operation = Operation.new(deadline: call_deadline) do |operation| + execute_request_response(method, request, marshal, unmarshal, metadata, credentials, operation) + end + return operation if return_op + + operation.execute + end + + # Close a channel created by this stub. + def close + @channel.close if @owned_channel + end + + private + + def execute_request_response(method, request, marshal, unmarshal, metadata, credentials, operation) + timeout = operation.deadline && operation.deadline - Time.now raise_deadline_exceeded if timeout && timeout <= 0 Sync do |task| if timeout - task.with_timeout(timeout) do - invoke_request_response(method, request, marshal, unmarshal, metadata, timeout) + task.with_timeout(timeout, Async::GRPC::DeadlineExceededError) do + metadata = update_metadata(metadata, credentials, method) + invoke_request_response(method, request, marshal, unmarshal, metadata, timeout, operation) end else - invoke_request_response(method, request, marshal, unmarshal, metadata, nil) + metadata = update_metadata(metadata, credentials, method) + invoke_request_response(method, request, marshal, unmarshal, metadata, nil, operation) end end - rescue Async::TimeoutError + rescue Async::GRPC::DeadlineExceededError raise_deadline_exceeded + rescue Async::GRPC::ResponseError => error + status = Protocol::GRPC::Status.for_http_status(error.response.status) + raise_bad_status(status, error.message, {}, cause: error) rescue Protocol::GRPC::Error => error raise_bad_status(error.status_code, error.cause&.message || error.message, error.metadata, cause: error) end - # Close a channel created by this stub. - def close - @channel.close if @owned_channel + def update_metadata(metadata, credentials, method) + metadata = normalize_metadata(metadata) + [@call_credentials, credentials].compact.each do |updater| + updater = updater.updater_proc if updater.respond_to?(:updater_proc) + raise TypeError, "Call credentials must be callable or expose updater_proc!" unless updater.respond_to?(:call) + + endpoint = @channel.endpoint + raise ArgumentError, "Call credentials require a secure channel with a known endpoint!" unless endpoint && endpoint.scheme == "https" + service = normalize_method(method).rpartition("/").first + context = {jwt_aud_uri: "https://#{endpoint.authority}#{service}"} + attributes = updater.call(context) + next if attributes.nil? + raise TypeError, "Call credentials must return a Hash or nil!" unless attributes.is_a?(Hash) + + # Google updaters can return the context along with authentication headers. + attributes.each do |key, value| + key = key.to_s + metadata[key] = value unless key == "jwt_aud_uri" + end + end + metadata end - private - - def invoke_request_response(method, request, marshal, unmarshal, metadata, timeout) + def invoke_request_response(method, request, marshal, unmarshal, metadata, timeout, operation) body = Protocol::GRPC::Body::Writable.new payload = marshal.call(request) - raise TypeError, "marshal must return a String" unless payload.is_a?(String) + raise TypeError, "Marshal must return a String!" unless payload.is_a?(String) body.write(payload) body.close_write + timeout = operation.deadline && operation.deadline - Time.now + raise_deadline_exceeded if timeout && timeout <= 0 + headers = build_headers( metadata: normalize_metadata(metadata), timeout: timeout, - content_type: "application/grpc+proto" + content_type: "application/grpc" ) request = Protocol::HTTP::Request["POST", normalize_method(method), headers, body] 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 = Protocol::GRPC::Body::Readable.wrap(response, encoding: response_encoding) payload = response_body&.read response_body&.finish + operation.trailing_metadata = extract_metadata(Protocol::HTTP::Headers.new(response.headers.trailer.to_a, policy: Protocol::GRPC::HEADER_POLICY)) + operation.status = ::Struct::Status.new( + Protocol::GRPC::Metadata.extract_status(response.headers), + Protocol::GRPC::Metadata.extract_message(response.headers), + operation.trailing_metadata + ) check_status!(response) payload ? unmarshal.call(payload) : nil @@ -244,19 +335,7 @@ def build_headers(metadata:, timeout:, content_type:) end def extract_metadata(headers) - metadata = {} - - headers.to_h.each do |key, value| - next if key.start_with?("grpc-") || key == "content-type" || key == "te" - - if key.end_with?("-bin") - value = value.map{|item| Base64.strict_decode64(item)} - end - - metadata[key] = value - end - - return metadata + Protocol::GRPC::Metadata.extract(headers) end def normalize_method(method) @@ -286,7 +365,7 @@ def relative_timeout(deadline) elsif deadline.is_a?(Numeric) deadline else - raise TypeError, "deadline must be a Time or Numeric value" + raise TypeError, "Deadline must be a Time or Numeric value!" end end @@ -297,11 +376,11 @@ def relative_default_timeout end def raise_deadline_exceeded - raise_bad_status(::GRPC::Core::StatusCodes::DEADLINE_EXCEEDED, "Deadline exceeded", {}) + raise_bad_status(::GRPC::Core::StatusCodes::DEADLINE_EXCEEDED, "Deadline exceeded!", {}) end def raise_bad_status(status, details, metadata, cause: nil) - error = ::GRPC::BadStatus.new_status_exception(status, details || "unknown cause", metadata) + error = ::GRPC::BadStatus.new_status_exception(status, details || "Unknown cause!", metadata) raise error, cause: cause end end diff --git a/lib/async/grpc/compatible/gapic.rb b/lib/async/grpc/compatible/gapic.rb new file mode 100644 index 0000000..0351a8c --- /dev/null +++ b/lib/async/grpc/compatible/gapic.rb @@ -0,0 +1,55 @@ +# frozen_string_literal: true + +# Released under the MIT License. +# Copyright, 2026, by Samuel Williams. + +require "gapic/grpc" +require_relative "client_stub" + +module Async + module GRPC + module Compatible + # Represents a GAPIC service stub that preserves Ruby credential updaters. + class GapicServiceStub < ::Gapic::ServiceStub + # Initialize a GAPIC adapter for a generated service definition. + # @parameter service [Class] The generated GRPC::GenericService definition. + # @parameter channel [Channel | Nil] An optional shared Async channel. + # @parameter options [Hash] GAPIC service options, including credentials and endpoint. + def initialize(service, channel: nil, **options) + @service_name = service.service_name + @async_channel = channel + super(ClientStub.for(service), **options) + end + + # Construct the Async stub before GAPIC converts credentials into opaque native objects. + # @parameter grpc_stub_class [Class] The compatible stub class. + # @parameter endpoint [String] The service endpoint. + # @parameter credentials [Object] The original GAPIC credentials. + # @parameter channel_args [Hash | Nil] Channel arguments. + # @parameter interceptors [Array | Nil] Client interceptors. + def create_grpc_stub(grpc_stub_class, endpoint:, credentials:, channel_args: nil, interceptors: nil) + @grpc_stub = grpc_stub_class.new(endpoint, credentials, + channel_override: @async_channel, + channel_args: channel_args || {}, + interceptors: interceptors || []) + end + + # Async::HTTP owns connection pooling; native GAPIC channel pools are unsupported. + 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. + def close + @grpc_stub&.close + end + + # Supply a service identity for the anonymous generated stub class. + # @private + def setup_logging(system_name: nil, service: nil, **options) + super(system_name: "async-grpc-compatible", service: @service_name, **options) + end + end + end + end +end diff --git a/lib/async/grpc/compatible/operation.rb b/lib/async/grpc/compatible/operation.rb new file mode 100644 index 0000000..385b7aa --- /dev/null +++ b/lib/async/grpc/compatible/operation.rb @@ -0,0 +1,95 @@ +# frozen_string_literal: true + +# Released under the MIT License. +# Copyright, 2026, by Samuel Williams. + +require "async" +require "grpc" + +module Async + module GRPC + module Compatible + # Represents a deferred unary call. Execute and cancel active calls in the same reactor. + class Operation + # Initialize a deferred call. + # @parameter deadline [Time | Nil] The absolute call deadline. + # @yields {|operation| ...} Executes the request and records its response. + def initialize(deadline: nil, &execute) + @deadline = deadline + @execute = execute + @mutex = Thread::Mutex.new + @executed = false + @finished = false + @cancelled = false + @task = nil + @thread = nil + @metadata = nil + @trailing_metadata = nil + @status = nil + end + + # @attribute [Time | Nil] The absolute call deadline. + attr_reader :deadline + # @attribute [Hash | Nil] The initial response metadata. + attr_accessor :metadata + # @attribute [Hash | Nil] The response trailers. + attr_accessor :trailing_metadata + # @attribute [Struct::Status | Nil] The completed call status. + attr_accessor :status + + # Execute this operation once, waiting for the response. + # @returns [Object] The decoded response. + # @raises [GRPC::BadStatus] If the RPC fails or is cancelled. + def execute + @mutex.synchronize do + raise RuntimeError, "Operation has already been executed!" if @executed + @executed = true + end + + begin + Sync do |parent| + task = @mutex.synchronize do + raise ::GRPC::Cancelled.new("Cancelled!") if @cancelled + @thread = Thread.current + @task = Async::Task.new(parent, finished: false){@execute.call(self)} + end + + task.run + result = task.wait + raise ::GRPC::Cancelled.new("Cancelled!") if cancelled? + result + ensure + task&.stop + end + rescue ::GRPC::BadStatus => error + @status = error.to_status + raise + ensure + @mutex.synchronize do + @finished = true + @task = nil + end + @execute = nil + end + end + + # Cancel a pending or active operation from its reactor. + def cancel + task = @mutex.synchronize do + return if @finished + raise ThreadError, "Cancel the operation from its reactor thread!" if @task && @thread != Thread.current + @cancelled = true + @task + end + task&.stop + end + + # Whether cancellation was requested or reported by the server. + # @returns [Boolean] Whether the operation was cancelled. + def cancelled? + @mutex.synchronize{@cancelled || @status&.code == ::GRPC::Core::StatusCodes::CANCELLED} + end + end + end + end +end diff --git a/readme.md b/readme.md index 1211a6b..ceb5cc5 100644 --- a/readme.md +++ b/readme.md @@ -4,7 +4,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 initial release still depends on `grpc` for its public credential and error types, but it does not use the native gRPC channel for requests. +The gem depends on `grpc` for service definitions and error types. TLS configuration comes from `IO::Endpoint`, and requests use Async's connection pool. ## Usage @@ -14,7 +14,7 @@ Select the compatible stub when constructing a generated client: require "async/grpc/compatible" stub_class = Async::GRPC::Compatible::ClientStub -stub = stub_class.new("grpc.example.com:443", GRPC::Core::ChannelCredentials.new) +stub = stub_class.new("grpc.example.com:443", IO::Endpoint::TLS::Configuration.new) response = stub.request_response( "/example.Service/Get", @@ -43,19 +43,125 @@ The initial implementation supports: - Unary `request_response` calls. - Custom marshal and unmarshal callables. - Request metadata and deadlines. - - Insecure and standard TLS endpoints. + - Insecure endpoints and TLS using `IO::Endpoint::TLS::Configuration`, including custom trust roots and client certificates, with `Compatible::ChannelCredentials.new` mapping gRPC's positional PEM arguments. - Translation of gRPC failures into `GRPC::BadStatus` subclasses. + - 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. The following are not yet supported: - - `return_op: true`. - Client, server, or bidirectional streaming. - grpc-ruby interceptors. - - Parent call propagation and per-call credentials. - - Custom TLS root certificates, client certificates, and native channel overrides. + - 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. - 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. + +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. + +Supply `call_credentials:` to the constructor for a default authentication callback, or `credentials:` to `request_response` for a per-call callback. Objects exposing `updater_proc`, such as Google authentication credentials, are also accepted. Callbacks run at execution time on every call, so token refreshes are used. + +Each callback receives a fresh authentication context containing `:jwt_aud_uri`, for example `https://grpc.example.com/example.Service`. The audience uses the actual channel's endpoint and RPC service path. Return a hash of authentication headers, or `nil` to add none. Returned headers are merged into a copy of the caller's metadata; per-call credentials run after constructor credentials. The audience context is never sent as a header. + +Authentication callbacks require a TLS channel with a known endpoint. A shared `Async::GRPC::Client` supplies its endpoint through its HTTP delegate; for a custom client, supply the endpoint explicitly when constructing `Compatible::Channel.new(endpoint, client: client)`. + +``` ruby +stub = Async::GRPC::Compatible::ClientStub.new( + "grpc.example.com:443", + IO::Endpoint::TLS::Configuration.new, + call_credentials: credentials.updater_proc +) +``` + +Passing an authentication callback as the constructor's second argument creates a default TLS channel with peer and hostname verification. For custom trust roots or mutual TLS, provide the TLS configuration separately: + +``` ruby +tls = IO::Endpoint::TLS::Configuration.new( + trust_store: IO::Endpoint::TLS::TrustStore.load("ca.pem"), + certificate_chain: IO::Endpoint::TLS::Certificates.parse(File.read("client-chain.pem")), + private_key: File.read("client-key.pem") +) + +stub = Async::GRPC::Compatible::ClientStub.new( + "grpc.example.com:443", tls, + call_credentials: credentials.updater_proc +) +``` + +TLS channels verify peers and hostnames by default, including localhost. An explicit URL must match the selected transport: TLS credentials require `https://`, and `:this_channel_is_insecure` requires `http://`. + +### gRPC TLS configuration + +`Async::GRPC::Compatible::ChannelCredentials.new` accepts the same three optional positional PEM arguments as `GRPC::Core::ChannelCredentials.new` and returns an `IO::Endpoint::TLS::Configuration` directly: + +| gRPC constructor argument | TLS configuration | +| --- | --- | +| Root certificates | `trust_store`, containing only the supplied roots | +| Client private key | `private_key` | +| Client certificate chain | `certificate_chain`, split into individual certificates in the supplied order | + +``` ruby +tls = Async::GRPC::Compatible::ChannelCredentials.new( + File.read("ca.pem"), + File.read("client-key.pem"), + File.read("client-chain.pem") +) + +stub = Async::GRPC::Compatible::ClientStub.new("grpc.example.com:443", tls) +``` + +Use `ChannelCredentials.new` for the transport's default trust store, or `ChannelCredentials.new(root_pem)` for custom roots without a client identity. The client key and certificate chain must be supplied together. Peer and hostname verification are always enabled by this mapping. gRPC-specific default-root overrides are not read; supply those roots explicitly. + +Pass the returned TLS configuration to `ClientStub` or `GapicServiceStub` at construction. Existing native credential objects cannot be converted through Ruby's public API, so retain the PEM inputs at that boundary. Supply Ruby authentication callbacks separately using `call_credentials:` on `ClientStub`. + +## GAPIC and generated Google clients + +Require the optional adapter after installing your Google client gem (which provides `gapic-common`). `GapicServiceStub` accepts the generated service definition and preserves the original credential updater before GAPIC wraps it in native credentials. Its `call_rpc` path retains GAPIC's retry policies and yields the completed operation to the caller. + +``` ruby +require "google/cloud/kms/v1" +require "google/cloud/kms/v1/service_services_pb" +require "async/grpc/compatible/gapic" + +credentials = Google::Auth.get_application_default( + ["https://www.googleapis.com/auth/cloud-platform"] +) + +Sync do + service = Async::GRPC::Compatible::GapicServiceStub.new( + Google::Cloud::Kms::V1::KeyManagementService::Service, + endpoint: "cloudkms.googleapis.com", + credentials: credentials + ) + + begin + request = Google::Cloud::Kms::V1::EncryptRequest.new( + name: "projects/my-project/locations/global/keyRings/my-ring/cryptoKeys/my-key", + plaintext: "hello" + ) + response = service.call_rpc(:encrypt, request, options: {timeout: 5}) do |response, operation| + puts operation.status.code + end + puts response.ciphertext.bytesize + ensure + service.close + end +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. + +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. + +For generated services without GAPIC, use `ClientStub.for(Service)` to construct a stub class with unary RPC methods. Streaming methods raise `NotImplementedError`. This adapter supports the operation methods listed above; native operation controls such as `start_call`, `wait`, and write flags are not implemented. + ## Development Run the test suite: diff --git a/releases.md b/releases.md index ad7a7c3..2f6da39 100644 --- a/releases.md +++ b/releases.md @@ -1,5 +1,15 @@ # Releases +## Unreleased + + - Send `application/grpc` and use shared metadata decoding, including unpadded binary metadata. + - Map invalid HTTP responses to grpc-ruby errors, preserving the native `ResponseError` and its buffered response as the cause. + - Support deferred unary operations with execution, cancellation, deadline, status, and response metadata access. + - Support Ruby authentication callbacks at stub construction and per call, evaluated on each execution with the service's JWT audience and merged into request metadata. + - 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. + ## v0.0.0 - Initial implementation of an Async-backed `GRPC::ClientStub` compatible unary client. diff --git a/test/async/grpc/compatible/channel_credentials.rb b/test/async/grpc/compatible/channel_credentials.rb new file mode 100644 index 0000000..68d9f6b --- /dev/null +++ b/test/async/grpc/compatible/channel_credentials.rb @@ -0,0 +1,67 @@ +# frozen_string_literal: true + +# Released under the MIT License. +# Copyright, 2026, by Samuel Williams. + +require "async/grpc/compatible/channel_credentials" +require "sus/fixtures/openssl" + +describe Async::GRPC::Compatible::ChannelCredentials do + include Sus::Fixtures::OpenSSL::ValidCertificateContext + + let(:roots) {certificate_authority_certificate.to_pem} + let(:chain) {certificate.to_pem + roots} + let(:private_key) {key.to_pem} + + it "uses the default trust store with peer verification when no roots are supplied" do + configuration = subject.new + + expect(configuration).to be_a(IO::Endpoint::TLS::Configuration) + expect(configuration.trust_store).to be_nil + expect(configuration.certificate_chain).to be_nil + expect(configuration.private_key).to be_nil + expect(configuration.verification).to be == :peer + end + + it "maps gRPC positional arguments and preserves certificate bundle ordering" do + configuration = subject.new(roots + certificate.to_pem, private_key, chain) + + expect(configuration.trust_store.certificates).to be == [roots.strip, certificate.to_pem.strip] + expect(configuration.trust_store.system_certificates?).to be == false + expect(configuration.certificate_chain).to be == [certificate.to_pem.strip, roots.strip] + expect(configuration.private_key).to be == private_key + expect(configuration.verification).to be == :peer + end + + it "allows a client identity with the default trust store" do + configuration = subject.new(nil, private_key, chain) + + expect(configuration.trust_store).to be_nil + expect(configuration.certificate_chain).to be == [certificate.to_pem.strip, roots.strip] + expect(configuration.private_key).to be == private_key + expect(configuration.verification).to be == :peer + end + + it "requires both the client certificate chain and private key" do + expect do + subject.new(roots, private_key) + end.to raise_exception(ArgumentError, message: be =~ /certificate chain and private key/) + + expect do + subject.new(roots, nil, chain) + end.to raise_exception(ArgumentError, message: be =~ /certificate chain and private key/) + end + + it "rejects empty or invalid certificate bundles" do + ["", "not a certificate"].each do |bundle| + expect{subject.new(bundle)}.to raise_exception(ArgumentError, message: be =~ /does not contain any certificates/) + expect{subject.new(roots, private_key, bundle)}.to raise_exception(ArgumentError, message: be =~ /does not contain any certificates/) + end + end + + it "rejects non-string TLS material" do + expect{subject.new(false)}.to raise_exception(TypeError) + expect{subject.new(roots, false, chain)}.to raise_exception(TypeError) + expect{subject.new(roots, private_key, false)}.to raise_exception(TypeError) + end +end diff --git a/test/async/grpc/compatible/client_stub.rb b/test/async/grpc/compatible/client_stub.rb index 90d84cb..432548a 100644 --- a/test/async/grpc/compatible/client_stub.rb +++ b/test/async/grpc/compatible/client_stub.rb @@ -4,12 +4,19 @@ # Copyright, 2026, by Samuel Williams. require "async/grpc/compatible" +require "async/grpc/compatible/gapic" require "async/grpc/dispatcher" require "async/grpc/service" require "base64" require "sus/fixtures/async/http" +require "googleauth" +require_relative "../../../fixtures/tls" class CompatibleMessage + def self.encode(message) + message.to_proto + end + def self.decode(payload) new(payload) end @@ -32,6 +39,14 @@ class CompatibleInterface < Protocol::GRPC::Interface streaming: :unary end +class GeneratedCompatibleService + include GRPC::GenericService + self.service_name = "compatible.Service" + self.marshal_class_method = :encode + self.unmarshal_class_method = :decode + rpc :Echo, CompatibleMessage, CompatibleMessage +end + class CompatibleService < Async::GRPC::Service def echo(input, output, call) request = input.read @@ -39,12 +54,16 @@ def echo(input, output, call) case request.value when "error" call.response.headers["x-error"] = "metadata" - call.response.headers["x-error-bin"] = Base64.strict_encode64("binary metadata") + call.response.headers["x-error-bin"] = Base64.strict_encode64("binary metadata").delete("=") Protocol::GRPC::Metadata.assign_status!( call.response.headers, status: Protocol::GRPC::Status::NOT_FOUND, message: "Missing" ) + when "content-type" + output.write(CompatibleMessage.new(call.request.headers["content-type"].to_s)) + when "auth" + output.write(CompatibleMessage.new(call.request.headers["authorization"].to_s)) when "slow" sleep(0.1) output.write(CompatibleMessage.new("slow")) @@ -79,7 +98,7 @@ def request(value, **options) end it "matches grpc-ruby's constructor call shape" do - compatible_parameters = subject.instance_method(:initialize).parameters + compatible_parameters = subject.instance_method(:initialize).parameters.reject{|type, name| name == :call_credentials} native_parameters = ::GRPC::ClientStub.instance_method(:initialize).parameters expect(compatible_parameters.map(&:first)).to be == native_parameters.map(&:first) @@ -101,6 +120,165 @@ def request(value, **options) expect(response.value).to be == "Hello" end + it "sends the grpc-ruby request content type" do + expect(request("content-type").value).to be == "application/grpc" + end + + with "credentials" do + include TLSContext + + it "updates call credentials each time without modifying caller metadata" do + count = 0 + updater = ->(metadata) do + count += 1 + metadata["authorization"] = "Bearer token-#{count}" + metadata + end + metadata = {"x-test" => "original"} + expect(request("auth", credentials: updater, metadata: metadata).value).to be == "Bearer token-1" + expect(request("auth", credentials: updater, metadata: metadata).value).to be == "Bearer token-2" + expect(metadata).to be == {"x-test" => "original"} + end + + it "supports credential objects at construction" do + credentials = Object.new + credentials.define_singleton_method(:updater_proc){->(metadata){metadata.merge("authorization" => "Bearer constructor")}} + credential_stub = subject.new("unused", credentials, channel_override: channel) + response = credential_stub.request_response("/#{service_name}/Echo", CompatibleMessage.new("auth"), CompatibleMessage.method(:encode), CompatibleMessage.method(:decode)) + expect(response.value).to be == "Bearer constructor" + end + + it "supports an explicit credential updater alongside TLS credentials" do + updater = ->(metadata){metadata.merge("authorization" => "Bearer explicit")} + credential_stub = subject.new("unused", tls_credentials, channel_override: channel, call_credentials: updater) + response = credential_stub.request_response("/#{service_name}/Echo", CompatibleMessage.new("auth"), CompatibleMessage.method(:encode), CompatibleMessage.method(:decode)) + expect(response.value).to be == "Bearer explicit" + end + + it "runs credential updaters when an operation executes" do + count = 0 + updater = ->(metadata){count += 1; metadata} + operation = request("Hello", credentials: updater, return_op: true) + expect(count).to be == 0 + operation.execute + expect(count).to be == 1 + end + + it "merges authentication headers without exposing request metadata to the callback" do + contexts = [] + updater = ->(context) do + contexts << context + {authorization: "Bearer token"} + end + metadata = {"x-test" => "original"}.freeze + + expect(request("Hello", credentials: updater, metadata: metadata).value).to be == "Hello:original" + expect(contexts).to be == [{jwt_aud_uri: "https://#{client_endpoint.authority}/#{service_name}"}] + expect(metadata).to be == {"x-test" => "original"} + end + + it "combines constructor and per-call authentication metadata" do + default = ->(context){{"x-test" => "default"}} + per_call = ->(context){{"x-test-bin" => "per-call"}} + credential_stub = subject.new("unused", nil, channel_override: channel, call_credentials: default) + response = credential_stub.request_response("/#{service_name}/Echo", CompatibleMessage.new("Hello"), CompatibleMessage.method(:encode), CompatibleMessage.method(:decode), credentials: per_call) + + expect(response.value).to be == "Hello:default:per-call" + end + + it "keeps request metadata when a callback returns nil" do + expect(request("Hello", credentials: ->(context){nil}, metadata: {"x-test" => "original"}).value).to be == "Hello:original" + end + + it "rejects callback results that are not metadata" do + expect(grpc_client).not.to receive(:call) + expect{request("Hello", credentials: ->(context){"Bearer token"})}.to raise_exception(TypeError, message: be == "Call credentials must return a Hash or nil!") + end + + it "does not send authentication context as request headers" do + headers = nil + mock(grpc_client) do |wrapper| + wrapper.wrap(:call) do |original, request| + headers = request.headers + original.call(request) + end + end + + expect(request("Hello", credentials: ->(context){context}).value).to be == "Hello" + expect(headers["jwt_aud_uri"]).to be_nil + end + end + + it "rejects call credentials on an insecure channel before invoking them" do + called = false + updater = ->(context){called = true; {authorization: "Bearer token"}} + expect(grpc_client).not.to receive(:call) + + expect{request("auth", credentials: updater)}.to raise_exception(ArgumentError, message: be =~ /secure channel/) + expect(called).to be == false + end + + it "rejects call credentials when a shared client's endpoint is unknown" do + unknown_channel = Async::GRPC::Compatible::Channel.new(client: Object.new) + unknown_stub = subject.new("example.googleapis.com", nil, channel_override: unknown_channel, call_credentials: ->(context){{}}) + + expect do + unknown_stub.request_response("/#{service_name}/Echo", CompatibleMessage.new("Hello"), CompatibleMessage.method(:encode), CompatibleMessage.method(:decode)) + end.to raise_exception(ArgumentError, message: be =~ /known endpoint/) + end + + it "rejects ambiguous constructor call credentials" do + updater = ->(context){{}} + expect{subject.new("example.googleapis.com", updater, call_credentials: updater)}.to raise_exception(ArgumentError, message: be == "Supply call credentials only once!") + end + + it "ignores cancellation after completion" do + operation = request("Hello", return_op: true) + operation.execute + operation.cancel + expect(operation).not.to be(:cancelled?) + expect(operation.status.code).to be == 0 + end + + it "preserves failed operation status and metadata" do + operation = request("error", return_op: true) + expect{operation.execute}.to raise_exception(::GRPC::NotFound) + expect(operation.status.code).to be == 5 + expect(operation.status.metadata["x-error-bin"]).to be == ["binary metadata"] + end + + it "can cancel an operation before execution" do + operation = request("Hello", return_op: true) + operation.cancel + expect{operation.execute}.to raise_exception(::GRPC::Cancelled) + expect(operation).to be(:cancelled?) + expect(operation.status.code).to be == 1 + end + + it "cancels an active operation without stopping its caller" do + operation = request("slow", return_op: true) + execution = Async do + expect{operation.execute}.to raise_exception(::GRPC::Cancelled) + :finished + end + Async::Task.current.sleep(0.01) + operation.cancel + expect(execution.wait).to be == :finished + expect(operation.status.code).to be == 1 + end + + it "counts time spent waiting to execute toward the deadline" do + operation = request("Hello", deadline: Time.now + 0.01, return_op: true) + Async::Task.current.sleep(0.02) + expect{operation.execute}.to raise_exception(::GRPC::DeadlineExceeded) + end + + it "preserves application decoder errors" do + expect do + stub.request_response("/#{service_name}/Echo", CompatibleMessage.new("Hello"), CompatibleMessage.method(:encode), ->(payload){raise IOError, "Application decoder!"}) + end.to raise_exception(IOError, message: be == "Application decoder!") + end + it "connects directly to a target" do direct_stub = subject.new(bound_url, :this_channel_is_insecure) response = direct_stub.request_response( @@ -222,7 +400,7 @@ def request(value, **options) it "rejects invalid deadlines" do expect do request("Hello", deadline: Object.new) - end.to raise_exception(TypeError, message: be =~ /deadline/) + end.to raise_exception(TypeError, message: be =~ /Deadline/) end it "rejects expired deadlines before making the request" do @@ -239,7 +417,7 @@ def request(value, **options) ->(_message){Object.new}, CompatibleMessage.method(:decode) ) - end.to raise_exception(TypeError, message: be =~ /marshal/) + end.to raise_exception(TypeError, message: be =~ /Marshal/) end it "translates protocol errors into grpc-ruby errors" do @@ -270,22 +448,24 @@ def request(value, **options) end end - it "rejects operation objects" do - expect do - request("Hello", return_op: true) - end.to raise_exception(NotImplementedError, message: be =~ /return_op/) + it "defers execution until the operation executes" do + operation = request("Hello", return_op: true) + expect(operation.status).to be_nil + expect(operation.execute.value).to be == "Hello" + expect(operation.status.code).to be == 0 + expect{operation.execute}.to raise_exception(RuntimeError, message: be =~ /already/) end it "rejects parent call propagation" do expect do request("Hello", parent: Object.new) - end.to raise_exception(NotImplementedError, message: be =~ /parent/) + end.to raise_exception(NotImplementedError, message: be =~ /Parent/) end it "rejects per-call credentials" do expect do request("Hello", credentials: Object.new) - end.to raise_exception(NotImplementedError, message: be =~ /credentials/) + end.to raise_exception(TypeError, message: be =~ /credentials/) end it "rejects interceptors" do @@ -294,6 +474,183 @@ def request(value, **options) end.to raise_exception(NotImplementedError, message: be =~ /interceptors/i) end + with "a non-gRPC upstream" do + let(:http_status) {503} + + let(:app) do + Protocol::HTTP::Middleware.for do |request| + Protocol::HTTP::Response[http_status, {"content-type" => "text/html", "x-request-id" => "123"}, ["", "Proxy failure!", ""]] + end + end + + {400 => 13, 401 => 16, 403 => 7, 404 => 12, 429 => 14, 502 => 14, 503 => 14, 504 => 14, 500 => 2, 200 => 2}.each do |http, grpc| + with "HTTP #{http}", http_status: http do + it "maps the status and preserves the response diagnostics" do + expect{request("Hello")}.to raise_exception(::GRPC::BadStatus, message: be(:include?, "Invalid gRPC response: HTTP #{http}, content-type \"text/html\"!")).and(have_attributes( + code: be == grpc, + cause: be_a(Async::GRPC::ResponseError).and(have_attributes( + response: have_attributes( + status: be == http, + headers: have_keys("x-request-id" => be == ["123"]), + body: be_a(Protocol::HTTP::Body::Buffered).and(have_attributes(join: be == "Proxy failure!")) + ) + )) + )) + end + end + end + end + + with "GAPIC" do + include TLSContext + + let(:updater) {->(metadata){metadata.merge("authorization" => "Bearer gapic")}} + let(:gapic) do + Async::GRPC::Compatible::GapicServiceStub.new(GeneratedCompatibleService, + endpoint: "example.googleapis.com", credentials: updater, channel: channel, logger: nil) + end + + it "uses the real GAPIC call path with original credentials and an operation" do + yielded = false + response = gapic.call_rpc(:echo, CompatibleMessage.new("auth")) do |response, operation| + yielded = true + expect(response.value).to be == "Bearer gapic" + expect(operation.status.code).to be == 0 + expect(operation.metadata).to be_a(Hash) + expect(operation.trailing_metadata).to be_a(Hash) + end + expect(response.value).to be == "Bearer gapic" + expect(yielded).to be == true + 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" + end + + it "rejects native GAPIC channel pooling" do + pool = Struct.new(:channel_count).new(2) + expect do + Async::GRPC::Compatible::GapicServiceStub.new(GeneratedCompatibleService, + endpoint: "example.googleapis.com", credentials: updater, channel_pool_config: pool, logger: nil) + end.to raise_exception(ArgumentError, message: be =~ /shared Async channel/) + end + + with "Google JWT credentials" do + let(:updater) do + Google::Auth::ServiceAccountJwtHeaderCredentials.new( + private_key: key.to_pem, + issuer: "unit@example.invalid", + project_id: "unit-project" + ) + end + + it "signs authentication for the actual shared channel's service audience" do + response = gapic.call_rpc(:echo, CompatibleMessage.new("auth")) + token = response.value.delete_prefix("Bearer ") + claims, header = JWT.decode(token, key.public_key, true, algorithm: "RS256", verify_aud: true, aud: "https://#{client_endpoint.authority}/#{service_name}") + + expect(claims["iss"]).to be == "unit@example.invalid" + expect(header["alg"]).to be == "RS256" + ensure + gapic.close + end + end + end + + with "TLS" do + include TLSContext + + it "maps gRPC root certificates and supports a separate authentication callback" do + credentials = Async::GRPC::Compatible::ChannelCredentials.new(certificate_authority_certificate.to_pem) + updater = ->(context){{"authorization" => "Bearer mapped"}} + direct_stub = subject.new(bound_url, credentials, call_credentials: updater) + response = direct_stub.request_response("/#{service_name}/Echo", CompatibleMessage.new("auth"), CompatibleMessage.method(:encode), CompatibleMessage.method(:decode)) + + expect(response.value).to be == "Bearer mapped" + ensure + direct_stub&.close + end + + it "uses custom trust roots for a direct TLS connection" do + direct_stub = subject.new(bound_url, tls_credentials) + response = direct_stub.request_response("/#{service_name}/Echo", CompatibleMessage.new("TLS"), CompatibleMessage.method(:encode), CompatibleMessage.method(:decode)) + + expect(response.value).to be == "TLS" + ensure + direct_stub&.close + end + + it "rejects a server outside the configured trust roots" do + direct_stub = subject.new(bound_url, IO::Endpoint::TLS::Configuration.new) + + 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/) + ensure + direct_stub&.close + end + + it "rejects a trusted certificate for the wrong hostname" do + wrong_endpoint = subject.endpoint_for("https://wrong.example.invalid", tls_credentials) + wrong_endpoint.endpoint = IO::Endpoint.tcp("127.0.0.1", client_endpoint.to_url.port) + wrong_channel = Async::GRPC::Compatible::Channel.new(wrong_endpoint) + wrong_stub = subject.new("unused", nil, channel_override: wrong_channel) + + 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/) + ensure + wrong_channel&.close + end + + with "mutual authentication" do + def server_tls_configuration + IO::Endpoint::TLS::Configuration.new( + trust_store: tls_credentials.trust_store, + certificate_chain: [certificate.to_pem], private_key: key.to_pem, + verification: :required + ) + end + + it "maps gRPC client credentials through GAPIC for mutual TLS" do + roots = certificate_authority_certificate.to_pem + credentials = Async::GRPC::Compatible::ChannelCredentials.new(roots, key.to_pem, certificate.to_pem + roots) + gapic = Async::GRPC::Compatible::GapicServiceStub.new(GeneratedCompatibleService, + endpoint: bound_url, credentials: credentials, logger: nil) + + expect(gapic.call_rpc(:echo, CompatibleMessage.new("mTLS")).value).to be == "mTLS" + ensure + gapic&.close + end + + it "presents the configured client certificate and private key" do + credentials = IO::Endpoint::TLS::Configuration.new( + trust_store: tls_credentials.trust_store, + certificate_chain: [certificate.to_pem, certificate_authority_certificate.to_pem], private_key: key.to_pem + ) + direct_stub = subject.new(bound_url, credentials) + response = direct_stub.request_response("/#{service_name}/Echo", CompatibleMessage.new("mTLS"), CompatibleMessage.method(:encode), CompatibleMessage.method(:decode)) + + expect(response.value).to be == "mTLS" + ensure + direct_stub&.close + end + + it "cannot connect without the required client certificate" do + direct_stub = subject.new(bound_url, tls_credentials) + + 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))) + ensure + direct_stub&.close + end + end + end + with ".setup_channel" do it "reuses a compatible channel" do expect(subject.setup_channel(channel, "unused", nil)).to be == channel @@ -303,13 +660,13 @@ def request(value, **options) compatible_channel = subject.setup_channel(grpc_client, "unused", nil) expect(compatible_channel.client).to be == grpc_client - expect(compatible_channel.endpoint).to be_nil + expect(compatible_channel.endpoint).to be == client_endpoint end it "rejects native channel overrides" do expect do subject.setup_channel(Object.new, "unused", nil) - end.to raise_exception(TypeError, message: be =~ /channel_override/) + end.to raise_exception(TypeError, message: be =~ /Channel override/) end end @@ -328,17 +685,66 @@ def request(value, **options) end it "constructs a secure HTTP/2 endpoint" do - credentials = ::GRPC::Core::ChannelCredentials.new + credentials = IO::Endpoint::TLS::Configuration.new endpoint = subject.endpoint_for("grpc.example.com:443", credentials) expect(endpoint.to_url.to_s).to be == "https://grpc.example.com/" expect(endpoint.protocol).to be == Async::HTTP::Protocol::HTTP2 end + it "verifies certificates and hostnames even on localhost" do + endpoint = subject.endpoint_for("localhost:443", IO::Endpoint::TLS::Configuration.new) + context = endpoint.endpoint.context + + expect(context.verify_mode).to be == OpenSSL::SSL::VERIFY_PEER + expect(context.verify_hostname).to be == true + expect(context.alpn_protocols).to be == ["h2"] + end + + it "verifies peers and hostnames with default compatible channel credentials" do + credentials = Async::GRPC::Compatible::ChannelCredentials.new + endpoint = subject.endpoint_for("localhost:443", credentials) + context = endpoint.endpoint.context + + expect(endpoint.scheme).to be == "https" + expect(context.verify_mode).to be == OpenSSL::SSL::VERIFY_PEER + expect(context.verify_hostname).to be == true + end + + it "rejects plaintext targets with compatible channel credentials" do + expect do + subject.endpoint_for("http://localhost:50051", Async::GRPC::Compatible::ChannelCredentials.new) + end.to raise_exception(ArgumentError, message: be =~ /scheme/) + end + + it "rejects plaintext targets with TLS credentials" do + expect do + subject.endpoint_for("http://example.googleapis.com", IO::Endpoint::TLS::Configuration.new) + end.to raise_exception(ArgumentError, message: be =~ /scheme/) + end + + it "rejects TLS targets with insecure credentials" do + expect{subject.endpoint_for("https://example.googleapis.com", :this_channel_is_insecure)}.to raise_exception(ArgumentError, message: be =~ /scheme/) + end + + it "rejects opaque native TLS credentials" do + expect do + subject.new("example.googleapis.com", ::GRPC::Core::ChannelCredentials.new) + end.to raise_exception(TypeError, message: be =~ /native gRPC credentials are unsupported/) + end + + it "rejects composed native credentials even when a shared channel is provided" do + credentials = ::GRPC::Core::ChannelCredentials.new.compose(::GRPC::Core::CallCredentials.new(->(context){{authorization: "Bearer token"}})) + + expect do + subject.new("example.googleapis.com", credentials, channel_override: channel) + end.to raise_exception(TypeError, message: be =~ /native gRPC credentials are unsupported/) + end + it "rejects invalid credentials" do expect do subject.endpoint_for("localhost:50051", nil) - end.to raise_exception(TypeError, message: be =~ /credentials/) + end.to raise_exception(TypeError, message: be =~ /Credentials/) end it "rejects unsupported target schemes" do diff --git a/test/async/grpc/compatible/gapic.rb b/test/async/grpc/compatible/gapic.rb new file mode 100644 index 0000000..c43f566 --- /dev/null +++ b/test/async/grpc/compatible/gapic.rb @@ -0,0 +1,57 @@ +# frozen_string_literal: true + +# Released under the MIT License. +# Copyright, 2026, by Samuel Williams. + +require "async/grpc/compatible/gapic" + +describe Async::GRPC::Compatible::GapicServiceStub do + let(:service) do + Class.new do + include ::GRPC::GenericService + self.service_name = "compatible.UnitService" + end + end + + let(:channel) {Async::GRPC::Compatible::Channel.new(client: Object.new)} + let(:credentials) {->(metadata){metadata.merge("authorization" => "Bearer token")}} + let(:stub) do + subject.new(service, endpoint: "example.googleapis.com", credentials: credentials, channel: channel, logger: nil) + end + + it "builds a compatible stub using the shared channel" do + expect(stub.grpc_stub).to be_a(Async::GRPC::Compatible::ClientStub) + expect(stub.grpc_stub.channel).to be(:equal?, channel) + expect(stub.channel_pool).to be_nil + end + + it "leaves the shared channel open when closed" do + expect(channel).not.to receive(:close) + + stub.close + 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/" + expect(owned_stub.grpc_stub.channel.client).to receive(:close) + + owned_stub.close + end + + it "rejects interceptors instead of silently discarding them" do + expect do + subject.new(service, endpoint: "example.googleapis.com", credentials: credentials, channel: channel, interceptors: [Object.new], logger: nil) + end.to raise_exception(NotImplementedError, message: be == "Client interceptors are not yet supported!") + end + + it "rejects native GAPIC channel pooling" do + pool = ::Gapic::ServiceStub::ChannelPool::Configuration.new + pool.channel_count = 2 + + expect do + subject.new(service, endpoint: "example.googleapis.com", credentials: credentials, channel_pool_config: pool, logger: nil) + end.to raise_exception(ArgumentError, message: be == "Use a shared Async channel instead of a GAPIC channel pool!") + end +end diff --git a/test/async/grpc/compatible/operation.rb b/test/async/grpc/compatible/operation.rb new file mode 100644 index 0000000..a4c2d3d --- /dev/null +++ b/test/async/grpc/compatible/operation.rb @@ -0,0 +1,182 @@ +# frozen_string_literal: true + +# Released under the MIT License. +# Copyright, 2026, by Samuel Williams. + +require "async/grpc/compatible/operation" +require "async/queue" +require "sus/fixtures/async" + +describe Async::GRPC::Compatible::Operation do + it "defers execution and yields itself to the call" do + calls = [] + deadline = Time.now + 60 + operation = subject.new(deadline: deadline) do |call| + calls << call + call.metadata = {"request-id" => "123"} + call.trailing_metadata = {"result" => "complete"} + :response + end + + expect(calls).to be == [] + expect(operation.deadline).to be == deadline + expect(operation.status).to be_nil + expect(operation).not.to be(:cancelled?) + expect(operation.execute).to be == :response + expect(calls).to be == [operation] + expect(operation.metadata).to be == {"request-id" => "123"} + expect(operation.trailing_metadata).to be == {"result" => "complete"} + + operation.cancel + expect(operation).not.to be(:cancelled?) + expect{operation.execute}.to raise_exception(RuntimeError, message: be == "Operation has already been executed!") + expect(calls).to be == [operation] + end + + it "does not execute a call cancelled before it starts" do + called = false + operation = subject.new{called = true} + operation.cancel + operation.cancel + + expect(operation).to be(:cancelled?) + expect{operation.execute}.to raise_exception(::GRPC::Cancelled) + expect(called).to be == false + expect(operation.status.code).to be == ::GRPC::Core::StatusCodes::CANCELLED + expect{operation.execute}.to raise_exception(RuntimeError, message: be == "Operation has already been executed!") + end + + it "preserves application errors and prevents executing the failed call again" do + error = IOError.new("Decoder failed!") + cleaned_up = false + operation = subject.new do + raise error + ensure + cleaned_up = true + end + + expect{operation.execute}.to raise_exception(IOError).and(be(:equal?, error)) + expect(cleaned_up).to be == true + expect(operation.status).to be_nil + operation.cancel + expect(operation).not.to be(:cancelled?) + expect{operation.execute}.to raise_exception(RuntimeError, message: be == "Operation has already been executed!") + end + + it "preserves gRPC errors and their status details and metadata" do + error = ::GRPC::NotFound.new("Missing!", {"reason" => "absent", "details-bin" => "\x00\xff".b}) + operation = subject.new{raise error} + + expect{operation.execute}.to raise_exception(::GRPC::NotFound).and(be(:equal?, error)) + expect(operation.status).to be == error.to_status + operation.cancel + expect(operation).not.to be(:cancelled?) + expect(operation.status).to be == error.to_status + end + + it "recognizes cancellation reported by the server" do + error = ::GRPC::Cancelled.new("Server cancelled!", {"reason" => "shutdown"}) + operation = subject.new{raise error} + + expect{operation.execute}.to raise_exception(::GRPC::Cancelled).and(be(:equal?, error)) + expect(operation).to be(:cancelled?) + expect(operation.status).to be == error.to_status + end + + with "an active call" do + include Sus::Fixtures::Async::ReactorContext + + let(:started) {Async::Queue.new} + let(:release) {Async::Queue.new} + let(:events) {[]} + let(:operation) do + subject.new do + started.enqueue(:started) + release.dequeue + ensure + events << :cleaned_up + end + end + + it "rejects repeated execution without disturbing the active call" do + execution = Async{operation.execute} + started.dequeue + + expect{operation.execute}.to raise_exception(RuntimeError, message: be == "Operation has already been executed!") + expect(events).to be == [] + release.enqueue(:response) + expect(execution.wait).to be == :response + expect(events).to be == [:cleaned_up] + end + + it "cancels the call and runs cleanup without stopping the caller" do + execution = Async do + operation.execute + rescue ::GRPC::Cancelled => error + events << :caller_continued + error + end + started.dequeue + + operation.cancel + expect(execution.wait).to be_a(::GRPC::Cancelled) + expect(events).to be == [:cleaned_up, :caller_continued] + expect(operation).to be(:cancelled?) + expect(operation.status.code).to be == ::GRPC::Core::StatusCodes::CANCELLED + operation.cancel + expect(events).to be == [:cleaned_up, :caller_continued] + end + + it "rejects execution from another thread and leaves the call cancellable" do + call = operation + execution = Async do + call.execute + rescue ::GRPC::Cancelled => error + error + end + started.dequeue + + error = Thread.new do + call.execute + rescue RuntimeError => error + error + end.value + + expect(error).to be_a(RuntimeError) + expect(error.message).to be == "Operation has already been executed!" + expect(events).to be == [] + call.cancel + expect(execution.wait).to be_a(::GRPC::Cancelled) + expect(events).to be == [:cleaned_up] + end + + it "rejects cancellation from another thread without cancelling the call" do + call = operation + execution = Async{call.execute} + started.dequeue + + error = Thread.new do + call.cancel + rescue ThreadError => error + error + end.value + + expect(error).to be_a(ThreadError) + expect(error.message).to be == "Cancel the operation from its reactor thread!" + expect(call).not.to be(:cancelled?) + expect(events).to be == [] + release.enqueue(:response) + expect(execution.wait).to be == :response + expect(events).to be == [:cleaned_up] + end + + it "cleans up the call when its caller is stopped" do + execution = Async{operation.execute} + started.dequeue + + execution.stop + expect(events).to be == [:cleaned_up] + expect{operation.execute}.to raise_exception(RuntimeError, message: be == "Operation has already been executed!") + end + end +end diff --git a/test/fixtures/tls.rb b/test/fixtures/tls.rb new file mode 100644 index 0000000..938c1ea --- /dev/null +++ b/test/fixtures/tls.rb @@ -0,0 +1,41 @@ +# frozen_string_literal: true + +# Released under the MIT License. +# Copyright, 2026, by Samuel Williams. + +require "sus/fixtures/openssl" + +module TLSContext + include Sus::Fixtures::OpenSSL::ValidCertificateContext + + def url + "https://127.0.0.1:0" + end + + def certificate + @tls_certificate ||= super.dup.tap do |certificate| + extensions = OpenSSL::X509::ExtensionFactory.new + certificate.add_extension(extensions.create_extension("subjectAltName", "DNS:localhost,IP:127.0.0.1,IP:::1")) + certificate.sign(certificate_authority_key, OpenSSL::Digest::SHA256.new) + end + end + + def tls_credentials + IO::Endpoint::TLS::Configuration.new( + trust_store: IO::Endpoint::TLS::TrustStore.parse(certificate_authority_certificate.to_pem) + ) + end + + def server_tls_configuration + IO::Endpoint::TLS::Configuration.new(certificate_chain: [certificate.to_pem], private_key: key.to_pem) + end + + def endpoint_options + super.merge(tls_configuration: server_tls_configuration) + end + + def make_client_endpoint(bound_endpoint) + port = bound_endpoint.sockets.first.to_io.local_address.ip_port + Async::GRPC::Compatible::ClientStub.endpoint_for("https://127.0.0.1:#{port}", tls_credentials) + end +end