aprs.me/lib/aprsme/is/is.ex
Graham McIntire 6eda1c0290
fix: resolve 19 bugs across Elixir, JS/TS, and config
Critical:
- Remove :protected from WeatherCache ETS (conflicted with :public)
- Register Aprsme.ReplayRegistry in supervision tree
- Route :continue_replay dead message to :start_replay handler
- Add 5000ms timeout to unbounded :rpc.call calls
- Fix falsy lat/lng check rejecting valid 0 coordinates

Medium:
- Fix live_reload pattern (temp_web -> aprsme_web)
- Remove duplicate ecto_repos config
- Make DB SSL configurable via DB_SSL env var
- Track sizeCheckTimeout on self for cleanup in destroyed()
- Wrap pushEvent calls in try/catch
- Remove unused size variable in marker cluster icon
- Capture touch coords at touchstart instead of stale TouchEvent
- Demote console.log to console.debug in production code
- Add catch-all to normalize_bounds preventing GenServer crash

Style:
- Add missing @impl true annotations in is.ex
- Improve migration failure error message
- Use textContent instead of innerHTML for error messages
- Use proper type for longPressTimer instead of any
2026-06-21 11:47:07 -05:00

720 lines
23 KiB
Elixir

defmodule Aprsme.Is do
@moduledoc false
use GenServer
alias Aprsme.Is.LoginParams
alias Aprsme.Is.PacketStats
require Logger
@aprs_timeout 60 * 1000
@keepalive_interval 20 * 1000
@backpressure_safety_valve_timeout 30_000
@fallback_server "noam.aprs2.net"
@fallback_timeout_seconds 300
defstruct [
:server,
:port,
:socket,
:timer,
:keepalive_timer,
:connected_at,
:packet_stats,
:login_params,
:safety_valve_timer,
:failure_started_at,
buffer: "",
backpressure_active: false
]
@type t :: %__MODULE__{
server: String.t(),
port: pos_integer(),
socket: :ssl.sslsocket() | nil,
timer: reference() | nil,
keepalive_timer: reference() | nil,
connected_at: DateTime.t() | nil,
packet_stats: PacketStats.t(),
buffer: String.t(),
login_params: LoginParams.t(),
backpressure_active: boolean(),
safety_valve_timer: reference() | nil,
failure_started_at: DateTime.t() | nil
}
def start_link(opts \\ []) do
GenServer.start_link(__MODULE__, opts, name: __MODULE__)
end
@impl true
def init(_opts) do
env = Application.get_env(:aprsme, :env)
disable_connection = Application.get_env(:aprsme, :disable_aprs_connection, false)
do_init(env, disable_connection)
end
defp do_init(:test, _) do
Logger.warning("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")
{:stop, :test_environment_disabled}
end
defp do_init(_env, false) do
# Trap exits so we can gracefully shut down
Process.flag(:trap_exit, true)
# Add a small delay to prevent rapid reconnection attempts
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
# 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.
server =
if Application.get_env(:aprsme, :aprs_is_using_fallback, false) do
Logger.warning("APRS-IS fallback active from previous failure, connecting to #{@fallback_server}")
@fallback_server
else
server
end
port = Application.get_env(:aprsme, :aprs_is_port, 14_580)
default_filter = Application.get_env(:aprsme, :aprs_is_default_filter, "r/33/-96/100")
aprs_user_id = Application.get_env(:aprsme, :aprs_is_login_id, "W5ISP")
aprs_passcode = Application.get_env(:aprsme, :aprs_is_password, "-1")
# Record connection start time
connected_at = DateTime.utc_now()
# Initialize packet statistics
packet_stats = %PacketStats{last_second_timestamp: System.system_time(:second)}
# Initialize state without requiring immediate connection
state = %__MODULE__{
server: server,
port: port,
connected_at: connected_at,
packet_stats: packet_stats,
login_params: %LoginParams{
user_id: aprs_user_id,
passcode: aprs_passcode,
filter: default_filter
},
failure_started_at: nil
}
# Try to connect initially, but don't fail if it doesn't work
case connect_to_aprs_is(server, port) do
{:ok, socket} ->
case send_login_string(socket, aprs_user_id, aprs_passcode, default_filter) do
:ok ->
timer = create_timer(@aprs_timeout)
keepalive_timer = create_keepalive_timer(@keepalive_interval)
{:ok, reset_failure(%{state | socket: socket, timer: timer, keepalive_timer: keepalive_timer})}
error ->
Logger.warning("Failed to login to APRS-IS: #{inspect(error)}, will retry")
:gen_tcp.close(socket)
schedule_reconnect(5000)
{:ok, track_failure(state)}
end
error ->
Logger.warning("Unable to establish initial connection to APRS-IS: #{inspect(error)}, will retry")
schedule_reconnect(5000)
{:ok, track_failure(state)}
end
end
# Client API
def stop do
Logger.info("Stopping Server")
GenServer.stop(__MODULE__, :stop)
end
def get_status do
status_from(Process.whereis(__MODULE__))
end
# GenServer isn't running — return a disconnected snapshot.
defp status_from(nil), do: disconnected_status(nil)
defp status_from(pid) when is_pid(pid) do
GenServer.call(__MODULE__, :get_status, Application.get_env(:aprsme, :is_call_timeout_ms, 5_000))
catch
# GenServer exists but isn't responding — treat as disconnected but
# fall back to the rotate.aprs2.net default server string.
:exit, _ -> disconnected_status(~c"rotate.aprs2.net")
end
# Build the "I'm not connected" status map. `default_server` is only used
# when the application config has no server set; pass nil to mean "use
# whatever is in config, or nil".
defp disconnected_status(default_server) do
server = Application.get_env(:aprsme, :aprs_is_server, default_server)
port = Application.get_env(:aprsme, :aprs_is_port, 14_580)
{stored_packet_count, oldest_packet_timestamp} = safe_packet_storage_stats()
%{
connected: false,
server: server_to_string(server),
port: port,
connected_at: nil,
uptime_seconds: 0,
login_id: Application.get_env(:aprsme, :aprs_is_login_id, "W5ISP"),
filter: Application.get_env(:aprsme, :aprs_is_default_filter, "r/33/-96/100"),
packet_stats: default_packet_stats(),
stored_packet_count: stored_packet_count,
oldest_packet_timestamp: oldest_packet_timestamp
}
end
defp safe_packet_storage_stats do
{Aprsme.Packets.get_total_packet_count(), Aprsme.Packets.get_oldest_packet_timestamp()}
rescue
DBConnection.OwnershipError ->
{0, nil}
catch
:exit, _ ->
{0, nil}
end
def set_filter(filter_string), do: send_message("#filter #{filter_string}")
def list_active_filters, do: send_message("#filter?")
def send_message(from, to, message) do
padded_callsign = String.pad_trailing(to, 9)
send_message("#{from}>APRS,TCPIP*::#{padded_callsign}:#{message}")
end
def send_message(message) do
case Process.whereis(__MODULE__) do
nil -> {:error, :not_connected}
_pid -> GenServer.call(__MODULE__, {:send_message, message})
end
end
# Server methods
@spec connect_to_aprs_is(String.t() | charlist(), pos_integer()) ::
{:ok, :ssl.sslsocket()} | {:error, any()}
defp connect_to_aprs_is(server, port) do
# Additional safeguard: prevent connections in test environment
env = Application.get_env(:aprsme, :env)
disable_connection = Application.get_env(:aprsme, :disable_aprs_connection, false)
do_connect_to_aprs_is(server, port, env, disable_connection)
end
defp do_connect_to_aprs_is(_server, _port, :test, _) do
Logger.warning("Attempted APRS-IS connection blocked in test environment")
{:error, :test_environment_blocked}
end
defp do_connect_to_aprs_is(_server, _port, _, true) do
Logger.warning("Attempted APRS-IS connection blocked in test environment")
{:error, :test_environment_blocked}
end
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
:gen_tcp.connect(String.to_charlist(server_string), port, opts)
end
@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
login_string =
"user #{aprs_user_id} pass #{aprs_passcode} vers aprs.me 0.1 filter #{filter}\r\n"
Logger.info("Sending login string: user #{aprs_user_id} pass ***** vers aprs.me 0.1 filter #{filter}")
:gen_tcp.send(socket, login_string)
end
@spec create_timer(non_neg_integer()) :: reference()
defp create_timer(timeout) do
Process.send_after(self(), :aprsme_no_message_timeout, timeout)
end
@spec create_keepalive_timer(non_neg_integer()) :: reference()
defp create_keepalive_timer(interval) do
Process.send_after(self(), :send_keepalive, interval)
end
@impl true
def handle_call({:send_message, message}, _from, state) do
case state.socket do
nil ->
Logger.warning("Attempted to send message while not connected to APRS-IS")
{:reply, {:error, :not_connected}, state}
socket ->
next_ack_number = :ets.update_counter(:aprsme, :message_number, 1)
# Append ack number
message = message <> "{" <> to_string(next_ack_number) <> "\r"
Logger.info("Sending message: #{inspect(message)}")
case :gen_tcp.send(socket, message) do
:ok ->
{:reply, :ok, state}
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)
:gen_tcp.close(socket)
schedule_reconnect(5000)
state = track_failure(%{state | socket: nil, timer: nil, keepalive_timer: nil})
{:reply, {:error, error}, state}
end
end
end
@impl true
def handle_call(:get_status, _from, state) do
connected = state.socket != nil
status = %{
connected: connected,
server: server_to_string(state.server),
port: state.port,
connected_at: if(connected, do: state.connected_at),
uptime_seconds: if(connected, do: DateTime.diff(DateTime.utc_now(), state.connected_at), else: 0),
login_id: state.login_params.user_id,
filter: state.login_params.filter,
packet_stats: state.packet_stats,
stored_packet_count: Aprsme.Packets.get_total_packet_count(),
oldest_packet_timestamp: Aprsme.Packets.get_oldest_packet_timestamp()
}
{:reply, status, state}
end
@impl true
def handle_info(:aprsme_no_message_timeout, state) do
case state.socket do
nil ->
Logger.debug("Timeout occurred while not connected - ignoring")
{:noreply, state}
_socket ->
Logger.error("Socket timeout detected. Killing genserver.")
{:stop, :aprsme_timeout, state}
end
end
@impl true
def handle_info(:send_keepalive, state) do
case state.socket do
nil ->
Logger.debug("Skipping keepalive - not connected to APRS-IS")
# Keep rescheduling so the timer stays alive for when we reconnect
keepalive_timer = create_keepalive_timer(@keepalive_interval)
{:noreply, %{state | keepalive_timer: keepalive_timer}}
socket ->
# Send a comment line as keepalive (APRS-IS standard)
case :gen_tcp.send(socket, "# keepalive\r\n") do
:ok ->
Logger.debug("Sent keepalive")
keepalive_timer = create_keepalive_timer(@keepalive_interval)
{:noreply, %{state | keepalive_timer: keepalive_timer}}
{:error, reason} ->
Logger.error("Failed to send keepalive: #{inspect(reason)}")
{:stop, :normal, state}
end
end
end
@impl true
def handle_info({:tcp, _socket, data}, state) do
handle_socket_data(data, state)
end
@impl true
def handle_info({:backpressure, :activate}, %{socket: nil} = state) do
{:noreply, state}
end
def handle_info({:backpressure, :activate}, %{backpressure_active: true} = state) do
{:noreply, state}
end
def handle_info({:backpressure, :activate}, state) do
:inet.setopts(state.socket, active: false)
if state.timer, do: Process.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")
{:noreply, %{state | backpressure_active: true, timer: nil, safety_valve_timer: safety_valve_timer}}
end
@impl true
def handle_info({:backpressure, :deactivate}, %{backpressure_active: false} = state) do
{:noreply, state}
end
def handle_info({:backpressure, :deactivate}, %{socket: nil} = state) do
cancel_safety_valve(state)
{:noreply, %{state | backpressure_active: false, safety_valve_timer: nil}}
end
def handle_info({:backpressure, :deactivate}, state) do
:inet.setopts(state.socket, active: true)
cancel_safety_valve(state)
timer = create_timer(@aprs_timeout)
Logger.info("Backpressure deactivated — socket resumed active mode")
{:noreply, %{state | backpressure_active: false, safety_valve_timer: nil, timer: timer}}
end
@impl true
def handle_info(:backpressure_safety_valve, %{backpressure_active: false} = state) do
{:noreply, state}
end
def handle_info(:backpressure_safety_valve, %{socket: nil} = state) do
{:noreply, %{state | backpressure_active: false, safety_valve_timer: nil}}
end
def handle_info(:backpressure_safety_valve, state) do
Logger.warning("Backpressure safety valve triggered — forcing resume after timeout")
:inet.setopts(state.socket, active: true)
timer = create_timer(@aprs_timeout)
{:noreply, %{state | backpressure_active: false, safety_valve_timer: nil, timer: timer}}
end
@impl true
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_safety_valve(state)
# Schedule reconnect
schedule_reconnect(5000)
state =
track_failure(%{
state
| socket: nil,
timer: nil,
keepalive_timer: nil,
backpressure_active: false,
safety_valve_timer: nil
})
{:noreply, state}
end
@impl true
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_safety_valve(state)
# Schedule reconnect
schedule_reconnect(5000)
state =
track_failure(%{
state
| socket: nil,
timer: nil,
keepalive_timer: nil,
backpressure_active: false,
safety_valve_timer: nil
})
{:noreply, state}
end
@impl true
def handle_info(:reconnect, state) do
Logger.info("Attempting to reconnect to APRS-IS...")
case connect_to_aprs_is(state.server, state.port) do
{:ok, socket} ->
case send_login_string(
socket,
state.login_params.user_id,
state.login_params.passcode,
state.login_params.filter
) do
:ok ->
Logger.info("Successfully reconnected to APRS-IS")
timer = create_timer(@aprs_timeout)
keepalive_timer = create_keepalive_timer(@keepalive_interval)
{:noreply,
reset_failure(%{
state
| socket: socket,
timer: timer,
keepalive_timer: keepalive_timer,
backpressure_active: false,
safety_valve_timer: nil
})}
error ->
Logger.warning("Failed to login to APRS-IS: #{inspect(error)}, will retry in 10 seconds")
:gen_tcp.close(socket)
schedule_reconnect(10_000)
{:noreply, track_failure(state)}
end
error ->
Logger.warning("Unable to reconnect to APRS-IS: #{inspect(error)}, will retry in 10 seconds")
schedule_reconnect(10_000)
{:noreply, track_failure(state)}
end
end
defp handle_socket_data(data, state) do
if state.timer, do: Process.cancel_timer(state.timer)
current_time = System.system_time(:second)
packet_stats = update_packet_stats(state.packet_stats, current_time)
buffer = state.buffer <> data
{complete_lines, remaining_buffer} = extract_complete_lines(buffer)
Enum.each(complete_lines, fn line ->
trimmed = String.trim(line)
if trimmed != "", do: dispatch(trimmed)
end)
timer = Process.send_after(self(), :aprsme_no_message_timeout, @aprs_timeout)
{:noreply, %{state | timer: timer, packet_stats: packet_stats, buffer: remaining_buffer}}
end
# Extract complete lines from buffer, returning {complete_lines, remaining_buffer}
@spec extract_complete_lines(String.t()) :: {[String.t()], String.t()}
defp extract_complete_lines(buffer) do
# Split by both \r\n and \n to handle different line endings
parts = String.split(buffer, ~r/\r?\n/, parts: :infinity)
# The last part might be incomplete
case parts do
[] ->
{[], ""}
[single] ->
# No newline found, entire buffer is incomplete
{[], single}
parts ->
# Last element might be incomplete line
{complete, [maybe_incomplete]} = Enum.split(parts, -1)
# If buffer ended with newline, maybe_incomplete will be empty string
if String.ends_with?(buffer, "\n") or String.ends_with?(buffer, "\r\n") do
{complete ++ [maybe_incomplete], ""}
else
{complete, maybe_incomplete}
end
end
end
@impl true
def terminate(reason, state) do
# Do Shutdown Stuff
Logger.info("Terminating APRS-IS connection: #{inspect(reason)}")
# Log any remaining buffered data
case Map.get(state, :buffer, "") do
"" -> :ok
buffer -> Logger.warning("Terminating with incomplete packet in buffer: #{inspect(buffer)}")
end
# Cancel timers
if timer = Map.get(state, :timer), do: Process.cancel_timer(timer)
if timer = Map.get(state, :keepalive_timer), do: Process.cancel_timer(timer)
if timer = Map.get(state, :safety_valve_timer), do: Process.cancel_timer(timer)
# Close socket
case Map.get(state, :socket) do
nil ->
:ok
socket ->
Logger.info("Closing socket")
:gen_tcp.close(socket)
end
:normal
end
@impl true
def code_change(_old_vsn, state, _extra) do
{:ok, state}
end
@spec dispatch(binary) :: nil | :ok
def dispatch("# logresp " <> rest) do
if String.contains?(rest, "unverified") do
Logger.warning("APRS-IS login unverified: #{String.trim(rest)}")
else
Logger.info("APRS-IS login accepted: #{String.trim(rest)}")
end
end
def dispatch("#" <> comment_text) do
Logger.debug("COMMENT: " <> String.trim(comment_text))
end
def dispatch(""), do: nil
def dispatch(message) do
case Aprs.parse(message) do
{:ok, parsed_message} ->
# Store the packet in the database for future replay
# Use GenStage pipeline for efficient batch processing
try do
# Set received_at and convert struct to map — PacketConsumer.prepare_packet_for_insert
# handles extract_additional_data and normalize_data_type, so don't duplicate here.
current_time = DateTime.truncate(DateTime.utc_now(), :microsecond)
attrs =
parsed_message
|> Map.put(:received_at, current_time)
|> Map.put(:raw, message)
|> struct_to_map()
Aprsme.PacketProducer.submit_packet(attrs)
rescue
error ->
Logger.error("Exception while submitting packet from #{inspect(parsed_message.sender)}: #{inspect(error)}")
Logger.debug("Raw message: #{inspect(message)}")
Logger.debug("Parsed message: #{inspect(parsed_message)}")
end
{:error, :invalid_packet} ->
Logger.debug("PARSE ERROR: invalid packet")
Aprsme.Packets.store_bad_packet(message, %{
message: "Invalid packet format",
type: "ParseError"
})
{:error, error} ->
Logger.debug("PARSE ERROR: " <> to_string(error))
Aprsme.Packets.store_bad_packet(message, %{message: error, type: "ParseError"})
end
end
@spec update_packet_stats(PacketStats.t(), integer()) :: PacketStats.t()
defp update_packet_stats(stats, current_time) do
new_total = stats.total_packets + 1
# Check if we need to reset the per-second counter
if current_time - stats.last_second_timestamp >= 1 do
%PacketStats{
total_packets: new_total,
last_packet_at: DateTime.utc_now(),
packets_per_second: 1,
last_second_count: 1,
last_second_timestamp: current_time
}
else
new_second_count = stats.last_second_count + 1
%{
stats
| total_packets: new_total,
last_packet_at: DateTime.utc_now(),
packets_per_second: new_second_count,
last_second_count: new_second_count
}
end
end
@spec server_to_string(String.t() | charlist() | any()) :: String.t()
defp server_to_string(server) when is_list(server), do: List.to_string(server)
defp server_to_string(server) when is_binary(server), do: server
defp server_to_string(server), do: to_string(server)
@spec default_packet_stats() :: PacketStats.t()
defp default_packet_stats do
%PacketStats{last_second_timestamp: System.system_time(:second)}
end
# Helper function to recursively convert structs to maps
# This handles nested structs that Map.from_struct/1 cannot handle
@spec struct_to_map(any()) :: any()
defp struct_to_map(%{__struct__: struct_type} = struct) do
converted_map =
struct
|> Map.from_struct()
|> Map.new(fn {k, v} -> {k, struct_to_map(v)} end)
# Add type information to help with later processing
Map.put(converted_map, :__original_struct__, struct_type)
end
defp struct_to_map(value) when is_list(value) do
Enum.map(value, &struct_to_map/1)
end
defp struct_to_map(value), do: value
defp cancel_safety_valve(%{safety_valve_timer: nil}), do: :ok
defp cancel_safety_valve(%{safety_valve_timer: timer}) do
Process.cancel_timer(timer)
:ok
end
@spec schedule_reconnect(non_neg_integer()) :: reference()
defp schedule_reconnect(delay) do
Process.send_after(self(), :reconnect, delay)
end
# Track a connection failure. If we've been failing continuously for
# @fallback_timeout_seconds, switch to the fallback server until next restart.
defp track_failure(state) do
now = DateTime.utc_now()
failure_started_at = state.failure_started_at || now
if DateTime.diff(now, failure_started_at) >= @fallback_timeout_seconds and
not Application.get_env(:aprsme, :aprs_is_using_fallback, false) do
Logger.warning(
"Failed to connect to APRS-IS for #{@fallback_timeout_seconds}s, " <>
"falling back to #{@fallback_server}"
)
Application.put_env(:aprsme, :aprs_is_using_fallback, true)
%{state | failure_started_at: failure_started_at, server: @fallback_server}
else
%{state | failure_started_at: failure_started_at}
end
end
# Reset the failure tracking on a successful connection.
defp reset_failure(state) do
%{state | failure_started_at: nil}
end
end