diff --git a/lib/aprsme/cluster/packet_distributor.ex b/lib/aprsme/cluster/packet_distributor.ex index a46c42a..a382428 100644 --- a/lib/aprsme/cluster/packet_distributor.ex +++ b/lib/aprsme/cluster/packet_distributor.ex @@ -16,16 +16,17 @@ defmodule Aprsme.Cluster.PacketDistributor do def distribute_packet(packet) do cluster_enabled = Application.get_env(:aprsme, :cluster_enabled, false) - - if cluster_enabled and LeaderElection.leader_cached?() do - Phoenix.PubSub.broadcast( - Aprsme.PubSub, - @pubsub_topic, - {:distributed_packet, packet} - ) - end + maybe_broadcast(cluster_enabled and LeaderElection.leader_cached?(), packet) end + # Only the leader broadcasts; non-leader clustered nodes and non-clustered + # nodes silently drop the packet. + defp maybe_broadcast(true, packet) do + Phoenix.PubSub.broadcast(Aprsme.PubSub, @pubsub_topic, {:distributed_packet, packet}) + end + + defp maybe_broadcast(false, _packet), do: :ok + def subscribe do Phoenix.PubSub.subscribe(Aprsme.PubSub, @pubsub_topic) end diff --git a/lib/aprsme/cluster/topology.ex b/lib/aprsme/cluster/topology.ex index f002a51..5d2374c 100644 --- a/lib/aprsme/cluster/topology.ex +++ b/lib/aprsme/cluster/topology.ex @@ -5,33 +5,33 @@ defmodule Aprsme.Cluster.Topology do require Logger def child_spec(opts) do - if Application.get_env(:aprsme, :cluster_enabled, false) do - topologies = Application.get_env(:libcluster, :topologies, []) + cluster_enabled = Application.get_env(:aprsme, :cluster_enabled, false) + topologies = Application.get_env(:libcluster, :topologies, []) + build_spec(cluster_enabled, topologies, opts) + end - # Log the configuration for debugging - Logger.info("Cluster.Topology starting with topologies: #{inspect(topologies)}") + # Clustering disabled — no-op spec. + defp build_spec(false, _topologies, opts), do: noop_spec(opts) - # Ensure we have valid topologies - if topologies == [] or topologies == nil do - Logger.warning("No libcluster topologies configured, clustering will not work") - # Return a no-op spec - %{ - id: __MODULE__, - start: {__MODULE__, :start_link, [opts]}, - type: :worker - } - else - supervisor_opts = Keyword.merge([name: Aprsme.ClusterSupervisor], opts) - {Cluster.Supervisor, [topologies, supervisor_opts]} - end - else - # Return a no-op spec when clustering is disabled - %{ - id: __MODULE__, - start: {__MODULE__, :start_link, [opts]}, - type: :worker - } - end + # Clustering enabled but no topologies configured — no-op + warning. + defp build_spec(true, topologies, opts) when topologies in [nil, []] do + Logger.warning("No libcluster topologies configured, clustering will not work") + noop_spec(opts) + end + + # Clustering enabled with topologies — delegate to libcluster's supervisor. + defp build_spec(true, topologies, opts) do + Logger.info("Cluster.Topology starting with topologies: #{inspect(topologies)}") + supervisor_opts = Keyword.merge([name: Aprsme.ClusterSupervisor], opts) + {Cluster.Supervisor, [topologies, supervisor_opts]} + end + + defp noop_spec(opts) do + %{ + id: __MODULE__, + start: {__MODULE__, :start_link, [opts]}, + type: :worker + } end def start_link(_opts), do: :ignore diff --git a/lib/aprsme_web/live/shared/packet_handler.ex b/lib/aprsme_web/live/shared/packet_handler.ex index e469467..9972e2f 100644 --- a/lib/aprsme_web/live/shared/packet_handler.ex +++ b/lib/aprsme_web/live/shared/packet_handler.ex @@ -21,22 +21,18 @@ defmodule AprsmeWeb.Live.SharedPacketHandler do process_fn = Keyword.fetch!(opts, :process_fn) enrich_packet? = Keyword.get(opts, :enrich_packet, true) - if filter_fn.(packet, socket) do - enriched_packet = - if enrich_packet? do - packet - |> sanitize_packet() - |> enrich_with_device_info() - else - packet - end - - process_fn.(enriched_packet, socket) - else - {:noreply, socket} - end + dispatch_packet_update(filter_fn.(packet, socket), packet, socket, process_fn, enrich_packet?) end + defp dispatch_packet_update(false, _packet, socket, _process_fn, _enrich?), do: {:noreply, socket} + + defp dispatch_packet_update(true, packet, socket, process_fn, enrich?) do + packet |> maybe_enrich(enrich?) |> process_fn.(socket) + end + + defp maybe_enrich(packet, true), do: packet |> sanitize_packet() |> enrich_with_device_info() + defp maybe_enrich(packet, false), do: packet + @doc """ Checks if packet sender matches the given callsign. """ @@ -109,13 +105,9 @@ defmodule AprsmeWeb.Live.SharedPacketHandler do Map.get(packet, :device_identifier) || Map.get(packet, "device_identifier") end - defp lookup_device_info(device_identifier) do - case device_identifier do - nil -> nil - "" -> nil - identifier -> DeviceCache.lookup_device(identifier) - end - end + defp lookup_device_info(nil), do: nil + defp lookup_device_info(""), do: nil + defp lookup_device_info(identifier), do: DeviceCache.lookup_device(identifier) defp nil_or_empty?(nil), do: true defp nil_or_empty?(""), do: true diff --git a/test/aprsme/cluster/topology_test.exs b/test/aprsme/cluster/topology_test.exs index 8674f32..3df18c4 100644 --- a/test/aprsme/cluster/topology_test.exs +++ b/test/aprsme/cluster/topology_test.exs @@ -47,5 +47,25 @@ defmodule Aprsme.Cluster.TopologyTest do assert topologies == [gossip: [strategy: Gossip]] assert Keyword.get(supervisor_opts, :name) == Aprsme.ClusterSupervisor end + + test "returns no-op spec when topologies is nil" do + Application.put_env(:aprsme, :cluster_enabled, true) + Application.put_env(:libcluster, :topologies, nil) + + spec = Topology.child_spec([]) + + assert %{id: Topology, start: {Topology, :start_link, [[]]}, type: :worker} = spec + end + + test "merges caller opts onto the default supervisor name" do + Application.put_env(:aprsme, :cluster_enabled, true) + Application.put_env(:libcluster, :topologies, gossip: [strategy: Gossip]) + + spec = Topology.child_spec(shutdown: :brutal_kill) + + assert {Cluster.Supervisor, [_, supervisor_opts]} = spec + assert supervisor_opts[:shutdown] == :brutal_kill + assert supervisor_opts[:name] == Aprsme.ClusterSupervisor + end end end diff --git a/test/aprsme_web/live/shared/packet_handler_test.exs b/test/aprsme_web/live/shared/packet_handler_test.exs new file mode 100644 index 0000000..1f53f19 --- /dev/null +++ b/test/aprsme_web/live/shared/packet_handler_test.exs @@ -0,0 +1,148 @@ +defmodule AprsmeWeb.Live.SharedPacketHandlerTest do + use ExUnit.Case, async: true + + alias AprsmeWeb.Live.SharedPacketHandler + + defp socket, do: %Phoenix.LiveView.Socket{assigns: %{__changed__: %{}}} + + defp allow_all, do: fn _packet, _socket -> true end + defp allow_none, do: fn _packet, _socket -> false end + + defp record_process(ref), + do: fn packet, s -> + send(ref, {:processed, packet}) + {:noreply, s} + end + + describe "handle_packet_update/3" do + test "skips processing when filter returns false" do + socket = socket() + {result, ^socket} = {nil, socket} + _ = result + + assert {:noreply, ^socket} = + SharedPacketHandler.handle_packet_update( + %{sender: "K5ABC"}, + socket, + filter_fn: allow_none(), + process_fn: record_process(self()) + ) + + refute_received {:processed, _} + end + + test "processes and enriches packet when filter returns true" do + {:noreply, _} = + SharedPacketHandler.handle_packet_update( + %{sender: "K5ABC"}, + socket(), + filter_fn: allow_all(), + process_fn: record_process(self()) + ) + + assert_received {:processed, packet} + # Enrichment added nil device_* keys even for unknown identifier. + assert Map.has_key?(packet, :device_model) + end + + test "skips enrichment when enrich_packet: false" do + {:noreply, _} = + SharedPacketHandler.handle_packet_update( + %{sender: "K5ABC"}, + socket(), + filter_fn: allow_all(), + process_fn: record_process(self()), + enrich_packet: false + ) + + assert_received {:processed, packet} + refute Map.has_key?(packet, :device_model) + end + end + + describe "packet_matches_callsign?/2" do + test "matches on atom :sender key" do + assert SharedPacketHandler.packet_matches_callsign?(%{sender: "K5ABC"}, "K5ABC") + end + + test "matches on string 'sender' key" do + assert SharedPacketHandler.packet_matches_callsign?(%{"sender" => "K5ABC"}, "K5ABC") + end + + test "string key wins over atom key when both present" do + # The implementation prefers the string key via Map.get fallback ordering. + assert SharedPacketHandler.packet_matches_callsign?( + %{"sender" => "K5ABC", sender: "OTHER"}, + "K5ABC" + ) + end + + test "comparison is case-insensitive via Callsign.matches?" do + assert SharedPacketHandler.packet_matches_callsign?(%{sender: "k5abc"}, "K5ABC") + end + + test "returns false when mismatched" do + refute SharedPacketHandler.packet_matches_callsign?(%{sender: "K5ABC"}, "W1XYZ") + end + end + + describe "has_weather_data?/1" do + test "returns true when at least one weather field is present" do + assert SharedPacketHandler.has_weather_data?(%{temperature: 72.0}) + assert SharedPacketHandler.has_weather_data?(%{"humidity" => 50}) + end + + test "returns false when no weather fields are present" do + refute SharedPacketHandler.has_weather_data?(%{sender: "K5ABC"}) + refute SharedPacketHandler.has_weather_data?(%{}) + end + end + + describe "enrich_with_device_info/1" do + test "adds nil device_* keys when no device identifier" do + result = SharedPacketHandler.enrich_with_device_info(%{sender: "K5ABC"}) + assert result.device_model == nil + assert result.device_vendor == nil + assert result.device_contact == nil + assert result.device_class == nil + end + + test "still adds nil device_* when device_identifier is empty string" do + result = SharedPacketHandler.enrich_with_device_info(%{device_identifier: ""}) + assert result.device_model == nil + end + end + + describe "enrich_packets_with_device_info/1" do + test "enriches each packet in the list" do + packets = [ + %{sender: "A", device_identifier: nil}, + %{sender: "B", device_identifier: ""}, + %{sender: "C", device_identifier: "APXYZ"} + ] + + enriched = SharedPacketHandler.enrich_packets_with_device_info(packets) + assert length(enriched) == 3 + assert Enum.all?(enriched, &Map.has_key?(&1, :device_model)) + end + + test "handles empty list" do + assert SharedPacketHandler.enrich_packets_with_device_info([]) == [] + end + end + + describe "filter factories" do + test "callsign_filter matches packets by sender" do + f = SharedPacketHandler.callsign_filter("K5ABC") + assert f.(%{sender: "K5ABC"}, :any) + refute f.(%{sender: "W1XYZ"}, :any) + end + + test "callsign_and_weather_filter requires both" do + f = SharedPacketHandler.callsign_and_weather_filter("K5ABC") + refute f.(%{sender: "K5ABC"}, :any) + refute f.(%{sender: "OTHER", temperature: 72.0}, :any) + assert f.(%{sender: "K5ABC", temperature: 72.0}, :any) + end + end +end