714 lines
23 KiB
Elixir
714 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
|
|
|
|
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
|
|
|
|
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
|
|
|
|
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
|
|
|
|
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
|
|
|
|
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
|
|
|
|
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
|
|
|
|
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
|
|
|
|
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 ->
|
|
require Logger
|
|
|
|
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
|