Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion async-grpc.gemspec
Original file line number Diff line number Diff line change
Expand Up @@ -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
19 changes: 16 additions & 3 deletions lib/async/grpc/client.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand All @@ -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

Expand All @@ -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
Expand Down
23 changes: 23 additions & 0 deletions lib/async/grpc/error.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
5 changes: 5 additions & 0 deletions releases.md
Original file line number Diff line number Diff line change
@@ -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.
Expand Down
2 changes: 1 addition & 1 deletion test/async/grpc/client.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
2 changes: 2 additions & 0 deletions test/async/grpc/dispatcher.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand All @@ -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

Expand Down
65 changes: 65 additions & 0 deletions test/async/grpc/interoperability.rb
Original file line number Diff line number Diff line change
@@ -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"}, ["<html>", "Proxy failure!", "</html>"]]
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 == "<html>Proxy failure!</html>"
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
Loading