From a25ee4f7632e97819b372cd1eda9ac1e37a2aae3 Mon Sep 17 00:00:00 2001 From: Michael Quinlan Date: Mon, 20 Jul 2026 14:05:37 -0600 Subject: [PATCH 1/2] fix: decode CreditResponse (0x8009) as non-correlated response The decoder listed :credit among correlated responses and tried to read a 4-byte correlation_id, but CreditResponse has none (code + subscription_id only), raising MatchError and crashing the connection. Decode it directly. --- lib/message/decoder.ex | 9 ++++++++- test/message/decoder_test.exs | 18 ++++++++++++++++++ 2 files changed, 26 insertions(+), 1 deletion(-) create mode 100644 test/message/decoder_test.exs 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..fd694b6 --- /dev/null +++ b/test/message/decoder_test.exs @@ -0,0 +1,18 @@ +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 +end From 085be44e815a6849905ba4e82947176cd2e00feb Mon Sep 17 00:00:00 2001 From: Michael Quinlan Date: Mon, 20 Jul 2026 14:06:27 -0600 Subject: [PATCH 2/2] fix: make parse_frames tolerant of undecodable frames Wrap per-frame Decoder.decode in try/rescue: log and drop a frame that fails to decode instead of crashing the connection GenServer. Defense in depth so no single unexpected frame can take down the stream pipeline. --- lib/message/buffer.ex | 14 +++++++++++++- test/message/decoder_test.exs | 15 +++++++++++++++ 2 files changed, 28 insertions(+), 1 deletion(-) 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/test/message/decoder_test.exs b/test/message/decoder_test.exs index fd694b6..7362a5e 100644 --- a/test/message/decoder_test.exs +++ b/test/message/decoder_test.exs @@ -15,4 +15,19 @@ defmodule RabbitMQStream.Message.DecoderTest do 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