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
93 changes: 70 additions & 23 deletions lib/prostaff_events/redis_subscriber.ex
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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)
Expand All @@ -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
Expand Down
63 changes: 63 additions & 0 deletions test/prostaff_events/redis_subscriber_test.exs
Original file line number Diff line number Diff line change
@@ -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
Expand Down Expand Up @@ -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
Loading