From cc88fd4a106568951be1b364ecbf2860b61f2c1f Mon Sep 17 00:00:00 2001 From: Graham McIntire Date: Sun, 26 Jul 2026 16:19:42 -0500 Subject: [PATCH] fix: dialyzer errors, pubsub failure handling, cache TTL, stream trim perf, and clearer warnings --- lib/aprsme/deployment_notifier.ex | 5 ++- lib/aprsme/device_cache.ex | 37 ++++++++++-------- lib/aprsme/is/is.ex | 39 ++++++++++++------- lib/aprsme_web/aprs_symbol.ex | 2 +- lib/aprsme_web/live/map_live/state.ex | 1 - lib/aprsme_web/live/map_live/subscriptions.ex | 21 ++++------ lib/aprsme_web/live/packets_live/index.ex | 34 ++++++++++------ lib/aprsme_web/live/status_live/index.ex | 2 +- 8 files changed, 83 insertions(+), 58 deletions(-) diff --git a/lib/aprsme/deployment_notifier.ex b/lib/aprsme/deployment_notifier.ex index a20ed14..8af4daf 100644 --- a/lib/aprsme/deployment_notifier.ex +++ b/lib/aprsme/deployment_notifier.ex @@ -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 diff --git a/lib/aprsme/device_cache.ex b/lib/aprsme/device_cache.ex index ea276a6..9ff1aad 100644 --- a/lib/aprsme/device_cache.ex +++ b/lib/aprsme/device_cache.ex @@ -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) diff --git a/lib/aprsme/is/is.ex b/lib/aprsme/is/is.ex index 41ab633..59f24cc 100644 --- a/lib/aprsme/is/is.ex +++ b/lib/aprsme/is/is.ex @@ -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) diff --git a/lib/aprsme_web/aprs_symbol.ex b/lib/aprsme_web/aprs_symbol.ex index ea6c126..844279c 100644 --- a/lib/aprsme_web/aprs_symbol.ex +++ b/lib/aprsme_web/aprs_symbol.ex @@ -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 diff --git a/lib/aprsme_web/live/map_live/state.ex b/lib/aprsme_web/live/map_live/state.ex index 6106868..bb44a10 100644 --- a/lib/aprsme_web/live/map_live/state.ex +++ b/lib/aprsme_web/live/map_live/state.ex @@ -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 """ diff --git a/lib/aprsme_web/live/map_live/subscriptions.ex b/lib/aprsme_web/live/map_live/subscriptions.ex index e1726d7..e64ee08 100644 --- a/lib/aprsme_web/live/map_live/subscriptions.ex +++ b/lib/aprsme_web/live/map_live/subscriptions.ex @@ -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 diff --git a/lib/aprsme_web/live/packets_live/index.ex b/lib/aprsme_web/live/packets_live/index.ex index 67138a5..02f6eec 100644 --- a/lib/aprsme_web/live/packets_live/index.ex +++ b/lib/aprsme_web/live/packets_live/index.ex @@ -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 """ diff --git a/lib/aprsme_web/live/status_live/index.ex b/lib/aprsme_web/live/status_live/index.ex index d7422e3..af3a5fd 100644 --- a/lib/aprsme_web/live/status_live/index.ex +++ b/lib/aprsme_web/live/status_live/index.ex @@ -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