diff --git a/async-grpc.gemspec b/async-grpc.gemspec index e68dabf..2ae88fd 100644 --- a/async-grpc.gemspec +++ b/async-grpc.gemspec @@ -28,6 +28,6 @@ Gem::Specification.new do |spec| spec.add_dependency "async", ">= 2.38.0" spec.add_dependency "async-http" - spec.add_dependency "protocol-grpc", "~> 0.14" + spec.add_dependency "protocol-grpc", "~> 0.16" spec.add_dependency "protocol-http", "~> 0.60" end diff --git a/lib/async/grpc/client.rb b/lib/async/grpc/client.rb index 80f3e98..eb831ed 100644 --- a/lib/async/grpc/client.rb +++ b/lib/async/grpc/client.rb @@ -98,11 +98,13 @@ def stub(interface_class, service_name) # Call the underlying HTTP client with merged headers. # @parameter request [Protocol::HTTP::Request] The HTTP request # @returns [Protocol::HTTP::Response] The HTTP response + # @raises [ResponseError] If the HTTP response does not conform to gRPC. def call(request) request.headers = @headers.merge(request.headers) super.tap do |response| response.headers.policy = Protocol::GRPC::HEADER_POLICY + validate_response!(response) end end @@ -117,16 +119,17 @@ def call(request) # @yields {|input, output| ...} Block for streaming calls # @returns [Object | Protocol::GRPC::Body::ReadableBody] Response message or readable body for streaming # @raises [ArgumentError] If method is unknown or streaming type is invalid + # @raises [ResponseError] If the HTTP response does not conform to gRPC. # @raises [Protocol::GRPC::Error] If the gRPC call fails def invoke(service, method, request = nil, metadata: {}, timeout: nil, encoding: nil, initial: nil, &block) rpc = service.class.lookup_rpc(method) - raise ArgumentError, "Unknown method: #{method}" unless rpc + raise ArgumentError, "Unknown method: #{method}!" unless rpc path = service.path(method) headers = Protocol::GRPC::Metadata.build( metadata: metadata, timeout: timeout, - content_type: "application/grpc+proto" + content_type: "application/grpc" ) headers["grpc-encoding"] = encoding if encoding @@ -144,12 +147,22 @@ def invoke(service, method, request = nil, metadata: {}, timeout: nil, encoding: when :bidirectional bidirectional_call(path, headers, request_class, response_class, encoding, initial: initial, &block) else - raise ArgumentError, "Unknown streaming type: #{streaming}" + raise ArgumentError, "Unknown streaming type: #{streaming}!" end end protected + # Reject non-gRPC responses before passing their bytes to a frame decoder. + # @parameter response [Protocol::HTTP::Response] The HTTP response. + # @raises [ResponseError] If the response is not a valid gRPC envelope. + def validate_response!(response) + content_type = response.headers["content-type"].to_s + return if response.status == 200 && content_type.match?(/\Aapplication\/grpc(?:\+[\w.-]+)?(?:\s*;|\z)/i) + + raise ResponseError.for(response) + end + # Make a unary gRPC call. # @parameter path [String] The gRPC path # @parameter headers [Protocol::HTTP::Headers] Request headers diff --git a/lib/async/grpc/error.rb b/lib/async/grpc/error.rb index dfd90e3..1c76ec8 100644 --- a/lib/async/grpc/error.rb +++ b/lib/async/grpc/error.rb @@ -14,6 +14,29 @@ class Error < StandardError class DeadlineExceededError < Error end + # Raised when an HTTP response does not conform to gRPC. + # Preserves the buffered response for inspection. + class ResponseError < Error + # Create an error from an invalid response, buffering its body for inspection. + # @parameter response [Protocol::HTTP::Response] The invalid response. + # @returns [ResponseError] The error with the buffered response attached. + def self.for(response) + self.new("Invalid gRPC response: HTTP #{response.status}, content-type #{response.headers["content-type"].to_s.inspect}!", response.buffered!) + end + + # Initialize an error with a message and response. + # @parameter message [String] The reason the response is invalid. + # @parameter response [Protocol::HTTP::Response] The buffered response. + def initialize(message, response) + super(message) + + @response = response + end + + # @attribute [Protocol::HTTP::Response] The buffered response, with its body available to read. + attr :response + end + # Represents an error that originated from a remote gRPC server. # Used as the `cause` of {Protocol::GRPC::Error} when the client receives a non-OK status. # The message and optional backtrace are extracted from response metadata. diff --git a/releases.md b/releases.md index 61e2c6a..97d13c8 100644 --- a/releases.md +++ b/releases.md @@ -1,5 +1,10 @@ # Releases +## Unreleased + + - Send `application/grpc` request headers for compatibility with Google API frontends. + - Reject invalid HTTP responses before decoding gRPC frames with `Async::GRPC::ResponseError`, describing the invalid status and content type in the message and buffering the response for inspecting status, headers, and body. + ## v0.9.0 - Added `Async::GRPC::Dispatcher#emit_completion` for once-per-request completion instrumentation, including routing failures, cancellations, and the final gRPC status. diff --git a/test/async/grpc/client.rb b/test/async/grpc/client.rb index a7337fa..97b717f 100644 --- a/test/async/grpc/client.rb +++ b/test/async/grpc/client.rb @@ -371,7 +371,7 @@ expect do client.invoke(interface, :InvalidCall) - end.to raise_exception(ArgumentError, message: be == "Unknown streaming type: invalid") + end.to raise_exception(ArgumentError, message: be == "Unknown streaming type: invalid!") end end diff --git a/test/async/grpc/dispatcher.rb b/test/async/grpc/dispatcher.rb index 7f7c2f7..20e6272 100644 --- a/test/async/grpc/dispatcher.rb +++ b/test/async/grpc/dispatcher.rb @@ -342,6 +342,7 @@ def close_streams(input, output, call) dispatcher = recording_dispatcher(error_service_name => error_service) response = dispatcher.call(build_request(error_service_name, "RaiseError")) + expect(response.headers["backtrace"]).to be_nil expect(Protocol::GRPC::Metadata.extract_status(response.headers)).to be == Protocol::GRPC::Status::RESOURCE_EXHAUSTED @@ -356,6 +357,7 @@ def close_streams(input, output, call) dispatcher = recording_dispatcher(error_service_name => error_service) response = dispatcher.call(build_request(error_service_name, "RaiseError")) + expect(response.headers["backtrace"]).to be_nil expect(Protocol::GRPC::Metadata.extract_status(response.headers)).to be == Protocol::GRPC::Status::INTERNAL diff --git a/test/async/grpc/interoperability.rb b/test/async/grpc/interoperability.rb new file mode 100644 index 0000000..d4900c3 --- /dev/null +++ b/test/async/grpc/interoperability.rb @@ -0,0 +1,65 @@ +# frozen_string_literal: true + +# Released under the MIT License. +# Copyright, 2026, by Samuel Williams. + +require "async/grpc/client" +require "protocol/http/body/writable" + +describe Async::GRPC::Client do + let(:request) {Protocol::HTTP::Request["POST", "/example.Service/Call", {}, nil]} + + def client_for(response) + delegate = Object.new + delegate.define_singleton_method(:call){|request| response} + subject.new(delegate) + end + + with "HTTP response validation" do + it "preserves the HTTP response and HTML body in the error" do + response = Protocol::HTTP::Response[503, {"content-type" => "text/html", "x-request-id" => "123"}, ["", "Proxy failure!", ""]] + expect do + client_for(response).call(request) + end.to raise_exception(Async::GRPC::ResponseError, message: be == "Invalid gRPC response: HTTP 503, content-type \"text/html\"!").and(have_attributes(response: be_equal(response))) + expect(response.headers["x-request-id"]).to be == ["123"] + expect(response.body).to be_a(Protocol::HTTP::Body::Buffered) + expect(response.read).to be == "Proxy failure!" + end + + it "rejects a non-200 status even with gRPC content type and OK status" do + response = Protocol::HTTP::Response[503, {"content-type" => "application/grpc", "grpc-status" => "0"}, nil] + expect do + client_for(response).call(request) + end.to raise_exception(Async::GRPC::ResponseError, message: be =~ /HTTP 503/) + expect(response.body).to be_nil + end + + it "closes the body when reading an invalid response fails" do + body = Protocol::HTTP::Body::Writable.new + body.define_singleton_method(:read){raise RuntimeError, "Read failed!"} + response = Protocol::HTTP::Response[503, {"content-type" => "text/html"}, body] + + expect{client_for(response).call(request)}.to raise_exception(RuntimeError, message: be == "Read failed!") + expect(body).to be(:closed?) + end + + ["application/grpc", "application/grpc+proto", "application/grpc+json", "application/grpc; charset=utf-8"].each do |content_type| + it "accepts #{content_type}" do + response = Protocol::HTTP::Response[200, {"content-type" => content_type, "grpc-status" => "0"}, ["unread"]] + expect(client_for(response).call(request)).to be_equal(response) + expect(response.body.read).to be == "unread" + ensure + response.close + end + end + + [nil, "application/grpc-web", "text/html"].each do |content_type| + it "rejects invalid content type #{content_type.inspect}" do + headers = {"grpc-status" => "0"} + headers["content-type"] = content_type if content_type + response = Protocol::HTTP::Response[200, headers, []] + expect{client_for(response).call(request)}.to raise_exception(Async::GRPC::ResponseError) + end + end + end +end