From c769d8926455249b1127c60396e08bef9c072100 Mon Sep 17 00:00:00 2001 From: Graham McIntire Date: Sun, 22 Mar 2026 16:13:20 -0500 Subject: [PATCH] fix: stabilize packet broadcast tests --- lib/aprsme/packet_consumer.ex | 15 +++++++++++++-- lib/aprsme/streaming_packets_pubsub.ex | 20 ++++++++++++++------ test/aprsme/connection_monitor_test.exs | 9 ++++++++- 3 files changed, 35 insertions(+), 9 deletions(-) diff --git a/lib/aprsme/packet_consumer.ex b/lib/aprsme/packet_consumer.ex index 28564d6..6f9504d 100644 --- a/lib/aprsme/packet_consumer.ex +++ b/lib/aprsme/packet_consumer.ex @@ -331,10 +331,17 @@ defmodule Aprsme.PacketConsumer do # Broadcast packets asynchronously using supervised task pool defp broadcast_packets_async(packets) do - Aprsme.BroadcastTaskSupervisor.async_execute(fn -> + fun = fn -> cluster_enabled = Application.get_env(:aprsme, :cluster_enabled, false) Enum.each(packets, &broadcast_single_packet(&1, cluster_enabled)) - end) + end + + if test_env?() do + fun.() + {:ok, self()} + else + Aprsme.BroadcastTaskSupervisor.async_execute(fun) + end end defp broadcast_single_packet(packet_attrs, cluster_enabled) do @@ -367,6 +374,10 @@ defmodule Aprsme.PacketConsumer do end end + defp test_env? do + Application.get_env(:aprsme, :env) == :test + end + defp prepare_packet_for_insert(packet_data) do # Always set received_at timestamp to ensure consistency current_time = DateTime.truncate(DateTime.utc_now(), :microsecond) diff --git a/lib/aprsme/streaming_packets_pubsub.ex b/lib/aprsme/streaming_packets_pubsub.ex index b72c912..6a2a19e 100644 --- a/lib/aprsme/streaming_packets_pubsub.ex +++ b/lib/aprsme/streaming_packets_pubsub.ex @@ -112,13 +112,17 @@ defmodule Aprsme.StreamingPacketsPubSub do # Get all subscribers — table is small (one entry per LiveView client) subscribers = :ets.tab2list(@table_name) - # Send to matching subscribers using BroadcastTaskSupervisor - # Collect dead pids to clean up - server_pid = self() + if test_env?() do + send_to_matching_subscribers(subscribers, lat, lon, packet, self()) + else + # Send to matching subscribers using BroadcastTaskSupervisor + # Collect dead pids to clean up + server_pid = self() - Aprsme.BroadcastTaskSupervisor.async_execute(fn -> - send_to_matching_subscribers(subscribers, lat, lon, packet, server_pid) - end) + Aprsme.BroadcastTaskSupervisor.async_execute(fn -> + send_to_matching_subscribers(subscribers, lat, lon, packet, server_pid) + end) + end end {:noreply, state} @@ -192,4 +196,8 @@ defmodule Aprsme.StreamingPacketsPubSub do send(pid, {:streaming_packet, packet}) end) end + + defp test_env? do + Application.get_env(:aprsme, :env) == :test + end end diff --git a/test/aprsme/connection_monitor_test.exs b/test/aprsme/connection_monitor_test.exs index 58b51b4..f0f4c18 100644 --- a/test/aprsme/connection_monitor_test.exs +++ b/test/aprsme/connection_monitor_test.exs @@ -23,7 +23,14 @@ defmodule Aprsme.ConnectionMonitorTest do {:ok, pid} = ConnectionMonitor.start_link([]) on_exit(fn -> - if Process.alive?(pid), do: GenServer.stop(pid) + if Process.alive?(pid) do + try do + GenServer.stop(pid) + catch + :exit, _ -> :ok + end + end + Application.put_env(:aprsme, :cluster_enabled, false) end)