fix: stabilize packet broadcast tests
This commit is contained in:
parent
ff81984a83
commit
c769d89264
3 changed files with 35 additions and 9 deletions
|
|
@ -331,10 +331,17 @@ defmodule Aprsme.PacketConsumer do
|
||||||
|
|
||||||
# Broadcast packets asynchronously using supervised task pool
|
# Broadcast packets asynchronously using supervised task pool
|
||||||
defp broadcast_packets_async(packets) do
|
defp broadcast_packets_async(packets) do
|
||||||
Aprsme.BroadcastTaskSupervisor.async_execute(fn ->
|
fun = fn ->
|
||||||
cluster_enabled = Application.get_env(:aprsme, :cluster_enabled, false)
|
cluster_enabled = Application.get_env(:aprsme, :cluster_enabled, false)
|
||||||
Enum.each(packets, &broadcast_single_packet(&1, cluster_enabled))
|
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
|
end
|
||||||
|
|
||||||
defp broadcast_single_packet(packet_attrs, cluster_enabled) do
|
defp broadcast_single_packet(packet_attrs, cluster_enabled) do
|
||||||
|
|
@ -367,6 +374,10 @@ defmodule Aprsme.PacketConsumer do
|
||||||
end
|
end
|
||||||
end
|
end
|
||||||
|
|
||||||
|
defp test_env? do
|
||||||
|
Application.get_env(:aprsme, :env) == :test
|
||||||
|
end
|
||||||
|
|
||||||
defp prepare_packet_for_insert(packet_data) do
|
defp prepare_packet_for_insert(packet_data) do
|
||||||
# Always set received_at timestamp to ensure consistency
|
# Always set received_at timestamp to ensure consistency
|
||||||
current_time = DateTime.truncate(DateTime.utc_now(), :microsecond)
|
current_time = DateTime.truncate(DateTime.utc_now(), :microsecond)
|
||||||
|
|
|
||||||
|
|
@ -112,13 +112,17 @@ defmodule Aprsme.StreamingPacketsPubSub do
|
||||||
# Get all subscribers — table is small (one entry per LiveView client)
|
# Get all subscribers — table is small (one entry per LiveView client)
|
||||||
subscribers = :ets.tab2list(@table_name)
|
subscribers = :ets.tab2list(@table_name)
|
||||||
|
|
||||||
# Send to matching subscribers using BroadcastTaskSupervisor
|
if test_env?() do
|
||||||
# Collect dead pids to clean up
|
send_to_matching_subscribers(subscribers, lat, lon, packet, self())
|
||||||
server_pid = self()
|
else
|
||||||
|
# Send to matching subscribers using BroadcastTaskSupervisor
|
||||||
|
# Collect dead pids to clean up
|
||||||
|
server_pid = self()
|
||||||
|
|
||||||
Aprsme.BroadcastTaskSupervisor.async_execute(fn ->
|
Aprsme.BroadcastTaskSupervisor.async_execute(fn ->
|
||||||
send_to_matching_subscribers(subscribers, lat, lon, packet, server_pid)
|
send_to_matching_subscribers(subscribers, lat, lon, packet, server_pid)
|
||||||
end)
|
end)
|
||||||
|
end
|
||||||
end
|
end
|
||||||
|
|
||||||
{:noreply, state}
|
{:noreply, state}
|
||||||
|
|
@ -192,4 +196,8 @@ defmodule Aprsme.StreamingPacketsPubSub do
|
||||||
send(pid, {:streaming_packet, packet})
|
send(pid, {:streaming_packet, packet})
|
||||||
end)
|
end)
|
||||||
end
|
end
|
||||||
|
|
||||||
|
defp test_env? do
|
||||||
|
Application.get_env(:aprsme, :env) == :test
|
||||||
|
end
|
||||||
end
|
end
|
||||||
|
|
|
||||||
|
|
@ -23,7 +23,14 @@ defmodule Aprsme.ConnectionMonitorTest do
|
||||||
{:ok, pid} = ConnectionMonitor.start_link([])
|
{:ok, pid} = ConnectionMonitor.start_link([])
|
||||||
|
|
||||||
on_exit(fn ->
|
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)
|
Application.put_env(:aprsme, :cluster_enabled, false)
|
||||||
end)
|
end)
|
||||||
|
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue