diff --git a/lib/prostaff_events/redis_subscriber.ex b/lib/prostaff_events/redis_subscriber.ex index b253e5d..0a03148 100644 --- a/lib/prostaff_events/redis_subscriber.ex +++ b/lib/prostaff_events/redis_subscriber.ex @@ -20,6 +20,15 @@ defmodule ProstaffEvents.RedisSubscriber do @redis_pattern "prostaff:events:*" + # Both the integer and the string form are accepted on purpose. + # Rails serializes `version` as a JSON number, so Jason.decode/1 hands us the + # integer 1. Comparing that against ["1"] never matches, and then + # resolve_topics/1 returns [] and the event reaches no topic at all - + # a silent drop with no error on either side. Do not narrow this list. + @supported_versions [1, "1"] + + @required_fields ["type", "org_id", "user_id"] + def start_link(_opts) do GenServer.start_link(__MODULE__, [], name: __MODULE__) end @@ -53,39 +62,74 @@ defmodule ProstaffEvents.RedisSubscriber do defp route_event(event) do case resolve_topics(event) do - [] -> - Logger.warning("[RedisSubscriber] Event missing required fields: #{inspect(event)}") + [] -> log_drop(event) + topics -> Enum.each(topics, &broadcast/1) + end + end + + defp broadcast({topic, message}) do + Phoenix.PubSub.broadcast(ProstaffEvents.PubSub, topic, message) + end - topics -> - Enum.each(topics, fn {topic, message} -> - Phoenix.PubSub.broadcast(ProstaffEvents.PubSub, topic, message) - end) + # resolve_topics/1 returns [] for two unrelated reasons, and calling both + # "missing required fields" points whoever is debugging at the wrong cause - + # the old message printed an event that had every required field in it. + # + # The payload is never logged. It can carry user data, and logs are retained + # and aggregated far more widely than the Redis channel they came from. + defp log_drop(event) do + case missing_fields(event) do + [] -> + Logger.warning( + "[RedisSubscriber] Rejected event: unsupported version=" <> + "#{inspect(field(event, "version"))} type=#{inspect(field(event, "type"))}" + ) + + missing -> + Logger.warning( + "[RedisSubscriber] Dropped event: missing or invalid field(s): " <> + Enum.join(missing, ", ") + ) end end - # Both the integer and the string form are accepted on purpose. - # Rails serializes `version` as a JSON number, so Jason.decode/1 hands us the - # integer 1. Comparing that against ["1"] never matches, and then - # resolve_topics/1 returns [] and the event reaches no topic at all - - # a silent drop with no error on either side. Do not narrow this list. - @supported_versions [1, "1"] + # A field only counts as present when it is a non-empty binary. A `type` that + # decodes as a number used to reach String.starts_with?/2 and raise inside + # handle_info, taking the subscriber down with it. + defp missing_fields(event) when is_map(event) do + Enum.reject(@required_fields, fn key -> + case Map.get(event, key) do + value when is_binary(value) -> value != "" + _ -> false + end + end) + end + + defp missing_fields(_event), do: @required_fields + + defp field(event, key) when is_map(event), do: Map.get(event, key) + defp field(_event, _key), do: nil @doc """ - Pure function: maps an event envelope to a list of {topic, message} pairs. - All broadcasting side effects happen outside this function, making it testable - without PubSub infrastructure. + Maps an event envelope to a list of {topic, message} pairs. - Rejects events with unknown versions. Accepts missing version for backward - compatibility with events published before versioning was introduced. + No broadcasting happens here, which is what makes the whole routing matrix + testable without Redis or PubSub. It does emit a warning for an event type + nothing routes - the type and org are only in scope here, and a silent drop + is how nine event types went unnoticed. Keep it that way. + + Returns [] for exactly two reasons: an unsupported version, or a required + field that is missing or not a binary. The caller tells them apart, so the + log names the right cause. + + A missing version is accepted, for compatibility with events published before + versioning existed. """ - def resolve_topics(%{"type" => type, "org_id" => org_id, "user_id" => user_id} = event) do + def resolve_topics(%{"type" => type, "org_id" => org_id, "user_id" => user_id} = event) + when is_binary(type) and is_binary(org_id) and is_binary(user_id) do version = Map.get(event, "version") if version != nil and version not in @supported_versions do - Logger.warning( - "[RedisSubscriber] Rejected event with unknown version=#{version} type=#{type}" - ) - [] else [{"org_events:#{org_id}", {:event, event}}] ++ specific_topics(type, org_id, user_id, event) @@ -106,7 +150,10 @@ defmodule ProstaffEvents.RedisSubscriber do [{"inhouse:#{org_id}", {:event, event}}] true -> - Logger.debug("[RedisSubscriber] Unrouted event type=#{type} org=#{org_id}") + # warning, not debug: production runs at :info and up, so an event type + # nobody routes was invisible. It means either a new type the publisher + # started sending or a typo - both need a human, and neither is normal. + Logger.warning("[RedisSubscriber] Unrouted event type=#{type} org=#{org_id}") [] end end diff --git a/test/prostaff_events/redis_subscriber_test.exs b/test/prostaff_events/redis_subscriber_test.exs index d08be0e..23bd488 100644 --- a/test/prostaff_events/redis_subscriber_test.exs +++ b/test/prostaff_events/redis_subscriber_test.exs @@ -1,6 +1,12 @@ defmodule ProstaffEvents.RedisSubscriberTest do use ExUnit.Case, async: true + import ExUnit.CaptureLog + + # US-08 subiu o ramo nao roteado para warning; sem isto ele polui a saida + # dos testes que nao estao olhando log. Falha ainda mostra tudo. + @moduletag :capture_log + alias ProstaffEvents.RedisSubscriber defp event(overrides \\ %{}) do @@ -105,4 +111,61 @@ defmodule ProstaffEvents.RedisSubscriberTest do assert RedisSubscriber.resolve_topics(e) == [] end end + + describe "motivo do descarte" do + test "versao nao suportada loga rejeicao por versao, nao campos ausentes" do + log = capture_log(fn -> send_event(event(%{"version" => 99})) end) + + assert log =~ "unsupported version" + refute log =~ "missing or invalid field" + end + + test "campo obrigatorio ausente loga qual campo faltou" do + log = capture_log(fn -> send_event(Map.delete(event(), "user_id")) end) + + assert log =~ "missing or invalid field" + assert log =~ "user_id" + refute log =~ "unsupported version" + end + + test "type nao-string e tratado como campo invalido, sem derrubar o processo" do + e = event(%{"type" => 123}) + + assert RedisSubscriber.resolve_topics(e) == [] + assert capture_log(fn -> send_event(e) end) =~ "type" + end + + test "nenhum dos logs de descarte imprime o payload" do + segredo = "nao-pode-vazar-no-log" + pid = start_subscriber() + + for e <- [ + event(%{"version" => 99, "payload" => %{"email" => segredo}}), + Map.delete(event(%{"payload" => %{"email" => segredo}}), "org_id") + ] do + refute capture_log(fn -> send_event(pid, e) end) =~ segredo + end + end + + test "tipo sem rota especifica loga em warning, nao em debug" do + log = capture_log(fn -> send_event(event(%{"type" => "roster.player_hired"})) end) + + assert log =~ "[warning]" + assert log =~ "Unrouted event" + assert log =~ "roster.player_hired" + end + end + + # route_event/1 e privado; o caminho publico e a mensagem do Redix. O + # subscriber registra com name: __MODULE__, entao e um por teste. + defp start_subscriber, do: start_supervised!({RedisSubscriber, []}, restart: :temporary) + + defp send_event(event), do: send_event(start_subscriber(), event) + + defp send_event(pid, event) do + send(pid, {:redix_pubsub, nil, nil, :pmessage, %{payload: Jason.encode!(event)}}) + # get_state serializa: quando retorna, o handle_info do evento ja terminou. + _ = :sys.get_state(pid) + :ok + end end