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
11 changes: 3 additions & 8 deletions lib/protocol/grpc/body/readable.rb
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,7 @@ def initialize(body, message_class: nil, encoding: nil)

# Read the next gRPC message.
# Overrides Wrapper#read to transform raw HTTP body chunks into decoded gRPC messages.
# Errors raised by the underlying body propagate unchanged.
# @returns [Object | String | Nil] Decoded message, raw binary, or `Nil` if stream ended
def read
# Read 5-byte prefix: 1 byte compression flag + 4 bytes length
Expand Down Expand Up @@ -85,14 +86,8 @@ def read
def read_exactly(n)
# Fill buffer until we have enough data:
while @buffer.bytesize < n
if @body.nil? || @body.empty?
return nil if @buffer.empty?

raise Error.new(Status::INTERNAL, "Truncated gRPC frame: expected #{n} bytes, received #{@buffer.bytesize}")
end

# Read chunk from underlying body:
chunk = @body.read
# An empty body can still have a pending error, so read to determine EOF:
chunk = @body&.read

if chunk.nil?
return nil if @buffer.empty?
Expand Down
4 changes: 4 additions & 0 deletions releases.md
Original file line number Diff line number Diff line change
@@ -1,5 +1,9 @@
# Releases

## Unreleased

- Preserve underlying body read errors, including failures before a message, during a partial frame, or while finishing the stream. These failures are no longer hidden as EOF or replaced with a truncated-frame error.

## v0.16.0

- Preserve timeout precision within the eight-digit wire limit, rounding up when necessary.
Expand Down
45 changes: 44 additions & 1 deletion test/protocol/grpc/body/readable.rb
Original file line number Diff line number Diff line change
@@ -1,10 +1,11 @@
# frozen_string_literal: true

# Released under the MIT License.
# Copyright, 2025, by Samuel Williams.
# Copyright, 2025-2026, by Samuel Williams.

require "protocol/grpc/body/readable"
require "protocol/http/body/buffered"
require "protocol/http/body/writable"
require_relative "../../../../fixtures/protocol/grpc/test_message"

require "zlib"
Expand Down Expand Up @@ -73,6 +74,48 @@ def write_data(data, compressed: false)
expect(body.read).to be_nil
end

it "returns nil when there is no underlying body" do
expect(subject.new(nil).read).to be_nil
end

with "pending body errors" do
let(:source_body) {Protocol::HTTP::Body::Writable.new}
let(:error) {IOError.new("Stream failed!")}

it "propagates an error even when the source body is empty" do
source_body.close_write(error)

expect(source_body).to be(:empty?)
expect{body.read}.to raise_exception(IOError).and(be_equal(error))
end

{
"partial prefix" => "\x00\x00".b,
"missing payload" => "\x00".b + [2].pack("N"),
"partial payload" => "\x00".b + [2].pack("N") + "a",
}.each do |description, chunk|
it "preserves the source error after a #{description}" do
source_body.write(chunk)
mock(source_body) do |wrapper|
wrapper.wrap(:read) do |original|
original.call.tap{source_body.close_write(error)}
end
end

expect{body.read}.to raise_exception(IOError).and(be_equal(error))
end
end

it "propagates an error while finishing after a complete message" do
message = message_class.new(value: "Hello")
write_message(message)
expect(body.read).to be == message
source_body.close_write(error)

expect{body.finish}.to raise_exception(IOError).and(be_equal(error))
end
end

it "works with binary mode (no message_class)" do
binary_body = subject.new(source_body, message_class: nil)
data = "Hello World".dup.force_encoding(Encoding::BINARY)
Expand Down
Loading