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
6 changes: 2 additions & 4 deletions lib/rocket/pusher.ex
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
29 changes: 17 additions & 12 deletions lib/rocket/request.ex
Original file line number Diff line number Diff line change
Expand Up @@ -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))
Expand Down
4 changes: 2 additions & 2 deletions mix.lock
Original file line number Diff line number Diff line change
Expand Up @@ -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"},
Expand Down
52 changes: 48 additions & 4 deletions test/rocket/pipeline_test.exs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
49 changes: 46 additions & 3 deletions test/rocket/request_test.exs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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"}}

Expand All @@ -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
Expand All @@ -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
Loading