diff --git a/lib/message/buffer.ex b/lib/message/buffer.ex index 3cb09f4..9082c66 100644 --- a/lib/message/buffer.ex +++ b/lib/message/buffer.ex @@ -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 diff --git a/lib/message/decoder.ex b/lib/message/decoder.ex index 3695d10..6e3e753 100644 --- a/lib/message/decoder.ex +++ b/lib/message/decoder.ex @@ -19,6 +19,14 @@ defmodule RabbitMQStream.Message.Decoder do |> decode(buffer) end + def decode(%Response{command: :credit} = response, buffer) do + <> = buffer + + response = %{response | code: decode_code(code)} + + %{response | data: Data.decode(response, "")} + end + def decode(%Response{command: command} = response, buffer) when command in [ :close, @@ -28,7 +36,6 @@ defmodule RabbitMQStream.Message.Decoder do :delete_producer, :subscribe, :unsubscribe, - :credit, :query_offset, :query_producer_sequence, :peer_properties, diff --git a/test/message/decoder_test.exs b/test/message/decoder_test.exs new file mode 100644 index 0000000..7362a5e --- /dev/null +++ b/test/message/decoder_test.exs @@ -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