Skip to content
Open
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
14 changes: 13 additions & 1 deletion lib/message/buffer.ex
Original file line number Diff line number Diff line change
Expand Up @@ -105,7 +105,19 @@ defmodule RabbitMQStream.Message.Buffer do
frames
|> Enum.reverse()
|> Enum.reduce(queue, fn frame, acc ->
:queue.in(Decoder.decode(frame), acc)
try do
:queue.in(Decoder.decode(frame), acc)
rescue
error ->
require Logger

Logger.warning("RabbitMQStream: dropping undecodable frame",
error: inspect(error),
frame: inspect(frame, limit: 64)
)

acc
end
end)
end
end
9 changes: 8 additions & 1 deletion lib/message/decoder.ex
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,14 @@ defmodule RabbitMQStream.Message.Decoder do
|> decode(buffer)
end

def decode(%Response{command: :credit} = response, buffer) do
<<code::unsigned-integer-size(16), _subscription_id::unsigned-integer-size(8)>> = buffer

response = %{response | code: decode_code(code)}

%{response | data: Data.decode(response, "")}
end

def decode(%Response{command: command} = response, buffer)
when command in [
:close,
Expand All @@ -28,7 +36,6 @@ defmodule RabbitMQStream.Message.Decoder do
:delete_producer,
:subscribe,
:unsubscribe,
:credit,
:query_offset,
:query_producer_sequence,
:peer_properties,
Expand Down
33 changes: 33 additions & 0 deletions test/message/decoder_test.exs
Original file line number Diff line number Diff line change
@@ -0,0 +1,33 @@
defmodule RabbitMQStream.Message.DecoderTest do
use ExUnit.Case, async: true

alias RabbitMQStream.Message.{Decoder, Response}
alias RabbitMQStream.Message.Types

# Exact prod CreditResponse frame (minus the 4-byte length prefix):
# key 0x8009, version 1, response_code 0x11 (precondition_failed), subscription_id 7.
@credit_response <<0x8009::16, 1::16, 0x11::16, 7::8>>

test "decodes a CreditResponse as a non-correlated response without crashing" do
result = Decoder.decode(@credit_response)

assert %Response{command: :credit, code: :precondition_failed} = result
assert result.correlation_id == nil
assert result.data == %Types.CreditResponseData{}
end

test "parse_frames drops an undecodable frame and keeps valid ones" do
alias RabbitMQStream.Message.Buffer

# 0x00FF is not a known command -> Decoder.decode raises. It must be skipped.
bad = <<0x00FF::16, 1::16, 0x00>>
good = @credit_response

# Buffer stores frames reversed (prepended); pass [good, bad] so processing
# order is bad then good.
queue = Buffer.parse_frames([good, bad], :queue.new())
commands = :queue.to_list(queue)

assert [%Response{command: :credit}] = commands
end
end