From 52f60ddd71934be7dc996a227c583c5dcb666783 Mon Sep 17 00:00:00 2001 From: Tom Pesman Date: Tue, 14 Jul 2026 09:39:39 +0200 Subject: [PATCH 1/2] Fix duplicate queued push error logging --- lib/rocket/pusher.ex | 6 ++-- lib/rocket/request.ex | 29 +++++++++++-------- test/rocket/pipeline_test.exs | 52 ++++++++++++++++++++++++++++++++--- test/rocket/request_test.exs | 49 +++++++++++++++++++++++++++++++-- 4 files changed, 113 insertions(+), 23 deletions(-) diff --git a/lib/rocket/pusher.ex b/lib/rocket/pusher.ex index d055c6e..673da7c 100644 --- a/lib/rocket/pusher.ex +++ b/lib/rocket/pusher.ex @@ -24,10 +24,8 @@ defmodule Rocket.Pusher do end defp perform_event(event) do - case request_module().perform(event) do - {:error, reason} -> Logger.error("[Rocket] push failed: #{inspect(reason)}") - _result -> :ok - end + request_module().perform(event) + :ok rescue error -> Logger.error("[Rocket] push raised: #{Exception.message(error)}") catch diff --git a/lib/rocket/request.ex b/lib/rocket/request.ex index 0eba936..7e24a0f 100644 --- a/lib/rocket/request.ex +++ b/lib/rocket/request.ex @@ -15,29 +15,34 @@ defmodule Rocket.Request do def perform(payload, opts) do handler = Keyword.get(opts, :response_handler, response_handler()) - with :ok <- validate_encodable(payload), - {:ok, response} <- post(payload) do + with {:ok, encoded_payload} <- encode(payload), + {:ok, response} <- post(encoded_payload) do handle_response(response, payload, handler) end end - defp validate_encodable(payload) do + defp encode(payload) do case Jason.encode(payload) do - {:ok, _encoded} -> :ok - {:error, error} -> {:error, {:encode_error, error}} + {:ok, encoded_payload} -> + {:ok, encoded_payload} + + {:error, error} -> + Logger.error("[Rocket] JSON encoding error #{inspect(error)}") + {:error, {:encode_error, error}} end end - defp post(payload) do - with {:ok, %{headers: headers, url: url}} <- config_provider().generate() do - {:ok, http_client().post(url, headers, Jason.encode!(payload), receive_timeout: 20_000)} + defp post(encoded_payload) do + case config_provider().generate() do + {:ok, %{headers: headers, url: url}} -> + {:ok, http_client().post(url, headers, encoded_payload, receive_timeout: 20_000)} + + {:error, reason} = error -> + Logger.error("[Rocket] configuration error #{inspect(reason)}") + error end - rescue - error in [Jason.EncodeError, Protocol.UndefinedError] -> {:error, {:encode_error, error}} end - defp handle_response({:error, {:encode_error, _error}} = error, _payload, _handler), do: error - defp handle_response({:ok, %{status: status}} = response, payload, handler) do parsed = Response.parse(response) handler.call(status, payload, response_body(parsed)) diff --git a/test/rocket/pipeline_test.exs b/test/rocket/pipeline_test.exs index 03a16c8..f5ddc24 100644 --- a/test/rocket/pipeline_test.exs +++ b/test/rocket/pipeline_test.exs @@ -2,6 +2,9 @@ defmodule Rocket.PipelineTest do use ExUnit.Case, async: false import ExUnit.CaptureLog + import Mox + + setup :verify_on_exit! defmodule TestConsumer do use GenStage @@ -30,7 +33,8 @@ defmodule Rocket.PipelineTest do case payload do :error -> {:error, :failed} :raise -> raise "request failed" - _payload -> :ok + :exit -> exit(:request_exited) + _payload -> {:ok, %{}} end end end @@ -146,17 +150,57 @@ defmodule Rocket.PipelineTest do test "pusher isolates failed events and continues processing" do Application.put_env(:rocket, :request_test_pid, self()) - capture_log(fn -> - assert {:noreply, [], :state} = Rocket.Pusher.handle_events([:error, :raise, :ok], self(), :state) - end) + log = + capture_log(fn -> + assert {:noreply, [], :state} = + Rocket.Pusher.handle_events([:error, :raise, :exit, :ok], self(), :state) + end) assert_receive {:performed, :error} assert_receive {:performed, :raise} + assert_receive {:performed, :exit} assert_receive {:performed, :ok} + + refute log =~ "[Rocket] push failed" + assert count_occurrences(log, "[Rocket] push raised: request failed") == 1 + assert count_occurrences(log, "[Rocket] push exited: {:exit, :request_exited}") == 1 + end + + test "pusher delegates an HTTP error to the configured response handler exactly once" do + Application.put_env(:rocket, :request_module, Rocket.Request) + Application.put_env(:rocket, :config_provider, Rocket.ConfigProviderMock) + Application.put_env(:rocket, :http_client, Rocket.HTTPClientMock) + Application.put_env(:rocket, :response_handler, Rocket.ResponseHandlerMock) + + payload = %{"message" => %{"token" => "bad-token"}} + + expect(Rocket.ConfigProviderMock, :generate, fn -> + {:ok, %{headers: [], url: "https://example.test/send"}} + end) + + expect(Rocket.HTTPClientMock, :post, fn _url, _headers, _body, _opts -> + {:ok, %{status: 400, body: ~s({"error":"invalid"})}} + end) + + expect(Rocket.ResponseHandlerMock, :call, fn 400, ^payload, %{"error" => "invalid"} -> :ok end) + + log = + capture_log(fn -> + assert {:noreply, [], :state} = Rocket.Pusher.handle_events([payload], self(), :state) + end) + + refute log =~ "[Rocket] push failed" end test "push collector callbacks keep producer state when there is no demand" do assert {:producer, :ok} = Rocket.PushCollector.init([]) assert {:noreply, [], :state} = Rocket.PushCollector.handle_demand(10, :state) end + + defp count_occurrences(log, message) do + log + |> String.split(message) + |> length() + |> Kernel.-(1) + end end diff --git a/test/rocket/request_test.exs b/test/rocket/request_test.exs index d940c35..1f2f5f9 100644 --- a/test/rocket/request_test.exs +++ b/test/rocket/request_test.exs @@ -6,6 +6,15 @@ defmodule Rocket.RequestTest do alias Rocket.Request + defmodule CustomResponseHandler do + @behaviour Rocket.Response.ResponseHandler + + @impl Rocket.Response.ResponseHandler + def call(status, payload, body) do + send(self(), {:custom_handler_called, status, payload, body}) + end + end + setup :verify_on_exit! setup do @@ -53,6 +62,23 @@ defmodule Rocket.RequestTest do assert Request.perform(payload) == {:error, %{"error" => "invalid"}} end + test "perform/2 preserves custom call/3 response handlers and structured results" do + payload = %{"message" => %{"token" => "bad-token"}} + + expect(Rocket.ConfigProviderMock, :generate, fn -> + {:ok, %{headers: [], url: "https://example.test/send"}} + end) + + expect(Rocket.HTTPClientMock, :post, fn _url, _headers, _body, _opts -> + {:ok, %{status: 400, body: ~s({"error":"invalid"})}} + end) + + assert Request.perform(payload, response_handler: CustomResponseHandler) == + {:error, %{"error" => "invalid"}} + + assert_receive {:custom_handler_called, 400, ^payload, %{"error" => "invalid"}} + end + test "returns invalid JSON errors from HTTP responses" do payload = %{"message" => %{"token" => "device-token"}} @@ -74,11 +100,21 @@ defmodule Rocket.RequestTest do test "returns configuration errors without posting" do expect(Rocket.ConfigProviderMock, :generate, fn -> {:error, :missing_credentials} end) - assert Request.perform(%{"message" => %{}}) == {:error, :missing_credentials} + log = + capture_log(fn -> + assert Request.perform(%{"message" => %{}}) == {:error, :missing_credentials} + end) + + assert count_occurrences(log, "[Rocket] configuration error :missing_credentials") == 1 end test "returns encode errors without configuration lookup" do - assert {:error, {:encode_error, %Protocol.UndefinedError{}}} = Request.perform(self()) + log = + capture_log(fn -> + assert {:error, {:encode_error, %Protocol.UndefinedError{}}} = Request.perform(self()) + end) + + assert count_occurrences(log, "[Rocket] JSON encoding error") == 1 end test "returns Finch transport errors and logs the connection issue" do @@ -97,6 +133,13 @@ defmodule Rocket.RequestTest do assert Request.perform(payload) == {:error, :timeout} end) - assert log =~ "[Rocket] connection error :timeout" + assert count_occurrences(log, "[Rocket] connection error :timeout") == 1 + end + + defp count_occurrences(log, message) do + log + |> String.split(message) + |> length() + |> Kernel.-(1) end end From d5fff06cc65e1bb023791a33f4573dc7ee143afc Mon Sep 17 00:00:00 2001 From: Tom Pesman Date: Tue, 14 Jul 2026 09:46:36 +0200 Subject: [PATCH 2/2] Update Mint security fixes --- mix.lock | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/mix.lock b/mix.lock index 88bdd42..d45ad82 100644 --- a/mix.lock +++ b/mix.lock @@ -9,14 +9,14 @@ "finch": {:hex, :finch, "0.19.0", "c644641491ea854fc5c1bbaef36bfc764e3f08e7185e1f084e35e0672241b76d", [:mix], [{:mime, "~> 1.0 or ~> 2.0", [hex: :mime, repo: "hexpm", optional: false]}, {:mint, "~> 1.6.2 or ~> 1.7", [hex: :mint, repo: "hexpm", optional: false]}, {:nimble_options, "~> 0.4 or ~> 1.0", [hex: :nimble_options, repo: "hexpm", optional: false]}, {:nimble_pool, "~> 1.1", [hex: :nimble_pool, repo: "hexpm", optional: false]}, {:telemetry, "~> 0.4 or ~> 1.0", [hex: :telemetry, repo: "hexpm", optional: false]}], "hexpm", "fc5324ce209125d1e2fa0fcd2634601c52a787aff1cd33ee833664a5af4ea2b6"}, "gen_stage": {:hex, :gen_stage, "1.2.1", "19d8b5e9a5996d813b8245338a28246307fd8b9c99d1237de199d21efc4c76a1", [:mix], [], "hexpm", "83e8be657fa05b992ffa6ac1e3af6d57aa50aace8f691fcf696ff02f8335b001"}, "goth": {:hex, :goth, "1.4.5", "ee37f96e3519bdecd603f20e7f10c758287088b6d77c0147cd5ee68cf224aade", [:mix], [{:finch, "~> 0.17", [hex: :finch, repo: "hexpm", optional: false]}, {:jason, "~> 1.1", [hex: :jason, repo: "hexpm", optional: false]}, {:jose, "~> 1.11", [hex: :jose, repo: "hexpm", optional: false]}], "hexpm", "0fc2dce5bd710651ed179053d0300ce3a5d36afbdde11e500d57f05f398d5ed5"}, - "hpax": {:hex, :hpax, "1.0.3", "ed67ef51ad4df91e75cc6a1494f851850c0bd98ebc0be6e81b026e765ee535aa", [:mix], [], "hexpm", "8eab6e1cfa8d5918c2ce4ba43588e894af35dbd8e91e6e55c817bca5847df34a"}, + "hpax": {:hex, :hpax, "1.0.4", "777de5d433b0fbdc7c418159c8055910faa8047ffdb3d6b31098d2a46cd7685c", [:mix], [], "hexpm", "afc7cb142ebcc2d01ce7816190b98ce5dd49e799111b24249f3443d730f377ca"}, "jason": {:hex, :jason, "1.4.5", "2e3a008590b0b8d7388c20293e9dcc9cf3e5d642fd2a114e4cbbb52e595d940a", [:mix], [{:decimal, "~> 1.0 or ~> 2.0 or ~> 3.0", [hex: :decimal, repo: "hexpm", optional: true]}], "hexpm", "b0c823996102bcd0239b3c2444eb00409b72f6a140c1950bc8b457d836b30684"}, "jose": {:hex, :jose, "1.11.10", "a903f5227417bd2a08c8a00a0cbcc458118be84480955e8d251297a425723f83", [:mix, :rebar3], [], "hexpm", "0d6cd36ff8ba174db29148fc112b5842186b68a90ce9fc2b3ec3afe76593e614"}, "makeup": {:hex, :makeup, "1.2.1", "e90ac1c65589ef354378def3ba19d401e739ee7ee06fb47f94c687016e3713d1", [:mix], [{:nimble_parsec, "~> 1.4", [hex: :nimble_parsec, repo: "hexpm", optional: false]}], "hexpm", "d36484867b0bae0fea568d10131197a4c2e47056a6fbe84922bf6ba71c8d17ce"}, "makeup_elixir": {:hex, :makeup_elixir, "1.0.1", "e928a4f984e795e41e3abd27bfc09f51db16ab8ba1aebdba2b3a575437efafc2", [:mix], [{:makeup, "~> 1.0", [hex: :makeup, repo: "hexpm", optional: false]}, {:nimble_parsec, "~> 1.2.3 or ~> 1.3", [hex: :nimble_parsec, repo: "hexpm", optional: false]}], "hexpm", "7284900d412a3e5cfd97fdaed4f5ed389b8f2b4cb49efc0eb3bd10e2febf9507"}, "makeup_erlang": {:hex, :makeup_erlang, "1.1.0", "835f7e60792e08824cda445639555d7bf1bbbddb1b60b306e33cb6f6db24dc74", [:mix], [{:makeup, "~> 1.0", [hex: :makeup, repo: "hexpm", optional: false]}], "hexpm", "1cd6780fb1dd1a03979abaed0fe82712b0625118fd5257d3ebbf73f960c73c3c"}, "mime": {:hex, :mime, "2.0.6", "8f18486773d9b15f95f4f4f1e39b710045fa1de891fada4516559967276e4dc2", [:mix], [], "hexpm", "c9945363a6b26d747389aac3643f8e0e09d30499a138ad64fe8fd1d13d9b153e"}, - "mint": {:hex, :mint, "1.7.1", "113fdb2b2f3b59e47c7955971854641c61f378549d73e829e1768de90fc1abf1", [:mix], [{:castore, "~> 0.1.0 or ~> 1.0", [hex: :castore, repo: "hexpm", optional: true]}, {:hpax, "~> 0.1.1 or ~> 0.2.0 or ~> 1.0", [hex: :hpax, repo: "hexpm", optional: false]}], "hexpm", "fceba0a4d0f24301ddee3024ae116df1c3f4bb7a563a731f45fdfeb9d39a231b"}, + "mint": {:hex, :mint, "1.9.2", "6e89e698d69cc29be001afd02c2cf4bae2d0994efe69a2ced7aab440d71584f9", [:mix], [{:castore, "~> 0.1.0 or ~> 1.0", [hex: :castore, repo: "hexpm", optional: true]}, {:hpax, "~> 0.1.1 or ~> 0.2.0 or ~> 1.0", [hex: :hpax, repo: "hexpm", optional: false]}], "hexpm", "d8e952b432fdac2321f570d29a68b9eb2664dc80a7a2ef9464531a5425abf222"}, "mix_audit": {:hex, :mix_audit, "2.1.5", "c0f77cee6b4ef9d97e37772359a187a166c7a1e0e08b50edf5bf6959dfe5a016", [:make, :mix], [{:jason, "~> 1.4", [hex: :jason, repo: "hexpm", optional: false]}, {:yaml_elixir, "~> 2.11", [hex: :yaml_elixir, repo: "hexpm", optional: false]}], "hexpm", "87f9298e21da32f697af535475860dc1d3617a010e0b418d2ec6142bc8b42d69"}, "mix_test_watch": {:hex, :mix_test_watch, "1.2.0", "1f9acd9e1104f62f280e30fc2243ae5e6d8ddc2f7f4dc9bceb454b9a41c82b42", [:mix], [{:file_system, "~> 0.2 or ~> 1.0", [hex: :file_system, repo: "hexpm", optional: false]}], "hexpm", "278dc955c20b3fb9a3168b5c2493c2e5cffad133548d307e0a50c7f2cfbf34f6"}, "mox": {:hex, :mox, "1.2.0", "a2cd96b4b80a3883e3100a221e8adc1b98e4c3a332a8fc434c39526babafd5b3", [:mix], [{:nimble_ownership, "~> 1.0", [hex: :nimble_ownership, repo: "hexpm", optional: false]}], "hexpm", "c7b92b3cc69ee24a7eeeaf944cd7be22013c52fcb580c1f33f50845ec821089a"},