fix: dialyzer errors, pubsub failure handling, cache TTL, stream trim perf, and clearer warnings
Some checks failed
Build and Push / Build and Push Docker Image (push) Failing after 2s
Some checks failed
Build and Push / Build and Push Docker Image (push) Failing after 2s
This commit is contained in:
parent
0e1161de39
commit
cc88fd4a10
8 changed files with 83 additions and 58 deletions
|
|
@ -14,7 +14,10 @@ defmodule Aprsme.DeploymentNotifier do
|
|||
|
||||
@impl true
|
||||
def init(_opts) do
|
||||
if System.get_env("DEPLOYED_AT"), do: Process.send_after(self(), :notify_deployment, 10_000)
|
||||
if System.get_env("DEPLOYED_AT"),
|
||||
do: :ok = Process.send_after(self(), :notify_deployment, 10_000),
|
||||
else: :ok
|
||||
|
||||
{:ok, %{}}
|
||||
end
|
||||
|
||||
|
|
|
|||
|
|
@ -97,26 +97,31 @@ defmodule Aprsme.DeviceCache do
|
|||
# Private functions
|
||||
|
||||
defp load_devices_into_cache do
|
||||
if Application.get_env(:aprsme, :env) == :test do
|
||||
Cache.put(@cache_name, :all_devices, [])
|
||||
:ok
|
||||
else
|
||||
devices =
|
||||
try do
|
||||
Repo.all(Devices)
|
||||
rescue
|
||||
error ->
|
||||
Logger.error("Failed to load devices from database: #{inspect(error)}")
|
||||
[]
|
||||
end
|
||||
devices = load_devices()
|
||||
|
||||
case Cache.put(@cache_name, :all_devices, devices) do
|
||||
{:ok, true} -> :ok
|
||||
error -> error
|
||||
end
|
||||
case Cache.put(@cache_name, :all_devices, devices) do
|
||||
{:ok, true} -> :ok
|
||||
_ -> :ok
|
||||
end
|
||||
end
|
||||
|
||||
defp load_devices, do: load_devices_impl()
|
||||
|
||||
defp load_devices_impl do
|
||||
case Application.get_env(:aprsme, :env) do
|
||||
:test -> []
|
||||
_ -> load_from_db()
|
||||
end
|
||||
end
|
||||
|
||||
defp load_from_db do
|
||||
Repo.all(Devices)
|
||||
rescue
|
||||
error ->
|
||||
Logger.error("Failed to load devices from database: #{inspect(error)}")
|
||||
[]
|
||||
end
|
||||
|
||||
defp find_matching_device(devices, identifier) do
|
||||
Enum.find(devices, fn device ->
|
||||
pattern_matches?(device.identifier, identifier)
|
||||
|
|
|
|||
|
|
@ -56,12 +56,12 @@ defmodule Aprsme.Is do
|
|||
end
|
||||
|
||||
defp do_init(:test, _) do
|
||||
Logger.warning("APRS-IS connection disabled in test environment")
|
||||
Logger.info("APRS-IS connection disabled in test environment")
|
||||
{:stop, :test_environment_disabled}
|
||||
end
|
||||
|
||||
defp do_init(_, true) do
|
||||
Logger.warning("APRS-IS connection disabled in test environment")
|
||||
Logger.warning("APRS-IS connection explicitly disabled")
|
||||
{:stop, :test_environment_disabled}
|
||||
end
|
||||
|
||||
|
|
@ -73,8 +73,7 @@ defmodule Aprsme.Is do
|
|||
Process.sleep(Application.get_env(:aprsme, :is_init_delay_ms, 2_000))
|
||||
|
||||
# Get startup parameters
|
||||
server_raw = Application.get_env(:aprsme, :aprs_is_server, ~c"dallas.aprs2.net")
|
||||
server = if is_list(server_raw), do: List.to_string(server_raw), else: server_raw
|
||||
server = initial_server()
|
||||
|
||||
# If we fell back to the fallback server in a previous failure streak and
|
||||
# the GenServer restarted (e.g. due to timeout), persist the fallback choice.
|
||||
|
|
@ -134,6 +133,11 @@ defmodule Aprsme.Is do
|
|||
end
|
||||
end
|
||||
|
||||
defp initial_server do
|
||||
server_raw = Application.get_env(:aprsme, :aprs_is_server, ~c"dallas.aprs2.net")
|
||||
to_server_string(server_raw)
|
||||
end
|
||||
|
||||
# Client API
|
||||
|
||||
def stop do
|
||||
|
|
@ -228,10 +232,13 @@ defmodule Aprsme.Is do
|
|||
defp do_connect_to_aprs_is(server, port, _env, false) do
|
||||
Logger.debug("Connecting to: #{server}:#{port}")
|
||||
opts = [:binary, active: true]
|
||||
server_string = if is_list(server), do: List.to_string(server), else: server
|
||||
server_string = to_server_string(server)
|
||||
:gen_tcp.connect(String.to_charlist(server_string), port, opts)
|
||||
end
|
||||
|
||||
defp to_server_string(server) when is_list(server), do: List.to_string(server)
|
||||
defp to_server_string(server) when is_binary(server), do: server
|
||||
|
||||
@spec send_login_string(:ssl.sslsocket(), String.t(), String.t(), String.t()) ::
|
||||
:ok | {:error, any()}
|
||||
defp send_login_string(socket, aprs_user_id, aprs_passcode, filter) do
|
||||
|
|
@ -273,8 +280,8 @@ defmodule Aprsme.Is do
|
|||
|
||||
error ->
|
||||
Logger.error("Failed to send message to APRS-IS: #{inspect(error)}")
|
||||
if state.timer, do: Process.cancel_timer(state.timer)
|
||||
if state.keepalive_timer, do: Process.cancel_timer(state.keepalive_timer)
|
||||
cancel_timer(state.timer)
|
||||
cancel_keepalive_timer(state.keepalive_timer)
|
||||
:gen_tcp.close(socket)
|
||||
schedule_reconnect(5000)
|
||||
|
||||
|
|
@ -358,7 +365,7 @@ defmodule Aprsme.Is do
|
|||
|
||||
def handle_info({:backpressure, :activate}, state) do
|
||||
:inet.setopts(state.socket, active: false)
|
||||
if state.timer, do: Process.cancel_timer(state.timer)
|
||||
cancel_timer(state.timer)
|
||||
safety_valve_timer = Process.send_after(self(), :backpressure_safety_valve, @backpressure_safety_valve_timeout)
|
||||
|
||||
Logger.warning("Backpressure activated — socket set to passive mode")
|
||||
|
|
@ -407,8 +414,8 @@ defmodule Aprsme.Is do
|
|||
def handle_info({:tcp_closed, _socket}, state) do
|
||||
Logger.warning("Socket has been closed by remote server - will reconnect")
|
||||
# Cancel any existing timers
|
||||
if state.timer, do: Process.cancel_timer(state.timer)
|
||||
if state.keepalive_timer, do: Process.cancel_timer(state.keepalive_timer)
|
||||
cancel_timer(state.timer)
|
||||
cancel_keepalive_timer(state.keepalive_timer)
|
||||
cancel_safety_valve(state)
|
||||
|
||||
# Schedule reconnect
|
||||
|
|
@ -431,8 +438,8 @@ defmodule Aprsme.Is do
|
|||
def handle_info({:tcp_error, _socket, reason}, state) do
|
||||
Logger.error("Connection error: #{inspect(reason)} - will reconnect")
|
||||
# Cancel any existing timers
|
||||
if state.timer, do: Process.cancel_timer(state.timer)
|
||||
if state.keepalive_timer, do: Process.cancel_timer(state.keepalive_timer)
|
||||
cancel_timer(state.timer)
|
||||
cancel_keepalive_timer(state.keepalive_timer)
|
||||
cancel_safety_valve(state)
|
||||
|
||||
# Schedule reconnect
|
||||
|
|
@ -493,7 +500,7 @@ defmodule Aprsme.Is do
|
|||
end
|
||||
|
||||
defp handle_socket_data(data, state) do
|
||||
if state.timer, do: Process.cancel_timer(state.timer)
|
||||
cancel_timer(state.timer)
|
||||
|
||||
current_time = System.system_time(:second)
|
||||
packet_stats = update_packet_stats(state.packet_stats, current_time)
|
||||
|
|
@ -688,6 +695,12 @@ defmodule Aprsme.Is do
|
|||
:ok
|
||||
end
|
||||
|
||||
defp cancel_timer(nil), do: :ok
|
||||
defp cancel_timer(ref), do: Process.cancel_timer(ref)
|
||||
|
||||
defp cancel_keepalive_timer(nil), do: :ok
|
||||
defp cancel_keepalive_timer(ref), do: Process.cancel_timer(ref)
|
||||
|
||||
@spec schedule_reconnect(non_neg_integer()) :: reference()
|
||||
defp schedule_reconnect(delay) do
|
||||
Process.send_after(self(), :reconnect, delay)
|
||||
|
|
|
|||
|
|
@ -184,7 +184,7 @@ defmodule AprsmeWeb.AprsSymbol do
|
|||
_ ->
|
||||
html = generate_marker_html(symbol_table, symbol_code, nil, size)
|
||||
# Cache for 1 hour since symbols don't change
|
||||
Aprsme.Cache.put(:symbol_cache, cache_key, html, ttl: Aprsme.Cache.to_timeout(hour: 1))
|
||||
_ = Aprsme.Cache.put(:symbol_cache, cache_key, html, ttl: Aprsme.Cache.to_timeout(hour: 1))
|
||||
html
|
||||
end
|
||||
else
|
||||
|
|
|
|||
|
|
@ -164,7 +164,6 @@ defmodule AprsmeWeb.MapLive.State do
|
|||
@doc """
|
||||
Default map zoom used by UrlParams.
|
||||
"""
|
||||
@spec default_zoom() :: integer()
|
||||
def default_zoom, do: UrlParams.default_zoom()
|
||||
|
||||
@doc """
|
||||
|
|
|
|||
|
|
@ -25,7 +25,7 @@ defmodule AprsmeWeb.MapLive.Subscriptions do
|
|||
def teardown(socket) do
|
||||
_teardown_connection_monitor(socket)
|
||||
_teardown_spatial(socket)
|
||||
_teardown_timers(socket)
|
||||
_teardown_all_timers(socket)
|
||||
_teardown_batch_tasks(socket)
|
||||
:ok
|
||||
end
|
||||
|
|
@ -99,20 +99,15 @@ defmodule AprsmeWeb.MapLive.Subscriptions do
|
|||
end
|
||||
end
|
||||
|
||||
defp _teardown_timers(socket) do
|
||||
_ = if socket.assigns.buffer_timer, do: Process.cancel_timer(socket.assigns.buffer_timer)
|
||||
|
||||
_ =
|
||||
if socket.assigns[:bounds_update_timer] do
|
||||
Process.cancel_timer(socket.assigns.bounds_update_timer)
|
||||
end
|
||||
|
||||
_ =
|
||||
if socket.assigns[:hover_end_timer] do
|
||||
Process.cancel_timer(socket.assigns.hover_end_timer)
|
||||
end
|
||||
defp _teardown_all_timers(socket) do
|
||||
cancel_if_exists(socket.assigns[:buffer_timer])
|
||||
cancel_if_exists(socket.assigns[:bounds_update_timer])
|
||||
cancel_if_exists(socket.assigns[:hover_end_timer])
|
||||
end
|
||||
|
||||
defp cancel_if_exists(nil), do: :ok
|
||||
defp cancel_if_exists(timer_ref), do: Process.cancel_timer(timer_ref)
|
||||
|
||||
defp _teardown_batch_tasks(socket) do
|
||||
_ =
|
||||
if socket.assigns[:pending_batch_tasks] do
|
||||
|
|
|
|||
|
|
@ -5,13 +5,21 @@ defmodule AprsmeWeb.PacketsLive.Index do
|
|||
|
||||
alias Aprsme.EncodingUtils
|
||||
|
||||
require Logger
|
||||
|
||||
@flush_interval_ms 150
|
||||
@max_stream_size 100
|
||||
|
||||
@impl true
|
||||
def mount(_params, _session, socket) do
|
||||
if connected?(socket) do
|
||||
Phoenix.PubSub.subscribe(Aprsme.PubSub, "postgres:aprsme_packets")
|
||||
case Phoenix.PubSub.subscribe(Aprsme.PubSub, "postgres:aprsme_packets") do
|
||||
:ok ->
|
||||
:ok
|
||||
|
||||
{:error, reason} ->
|
||||
Logger.error("Failed to subscribe to postgres:aprsme_packets: #{inspect(reason)}")
|
||||
end
|
||||
end
|
||||
|
||||
{:ok,
|
||||
|
|
@ -62,19 +70,21 @@ defmodule AprsmeWeb.PacketsLive.Index do
|
|||
|
||||
defp trim_stream(socket) do
|
||||
stream = socket.assigns.streams.packets
|
||||
count = Enum.count(stream)
|
||||
|
||||
if count > @max_stream_size do
|
||||
excess = count - @max_stream_size
|
||||
# Remove excess from the end (oldest items)
|
||||
to_remove = stream |> Enum.reverse() |> Enum.take(excess)
|
||||
stream
|
||||
|> Enum.to_list()
|
||||
|> Enum.reverse()
|
||||
|> do_trim(socket)
|
||||
end
|
||||
|
||||
Enum.reduce(to_remove, socket, fn {id, _}, s ->
|
||||
stream_delete(s, :packets, id)
|
||||
end)
|
||||
else
|
||||
socket
|
||||
end
|
||||
defp do_trim(packets, socket) when length(packets) <= @max_stream_size, do: socket
|
||||
|
||||
defp do_trim(packets, socket) do
|
||||
{_kept, to_remove} = Enum.split(packets, @max_stream_size)
|
||||
|
||||
Enum.reduce(to_remove, socket, fn {id, _}, s ->
|
||||
stream_delete(s, :packets, id)
|
||||
end)
|
||||
end
|
||||
|
||||
@doc """
|
||||
|
|
|
|||
|
|
@ -422,7 +422,7 @@ defmodule AprsmeWeb.StatusLive.Index do
|
|||
_ ->
|
||||
# Fallback to direct query if cache miss
|
||||
status = get_aprs_status()
|
||||
Aprsme.Cache.put(:query_cache, "aprs_status", status)
|
||||
Aprsme.Cache.put(:query_cache, "aprs_status", status, ttl: Aprsme.Cache.to_timeout(second: 10))
|
||||
status
|
||||
end
|
||||
end
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue