buffer incoming packets
This commit is contained in:
parent
4d2dedd751
commit
2dd5571577
4 changed files with 124 additions and 72 deletions
|
|
@ -60,6 +60,7 @@ defmodule Aprs.Is do
|
|||
keepalive_timer: keepalive_timer,
|
||||
connected_at: connected_at,
|
||||
packet_stats: packet_stats,
|
||||
buffer: "",
|
||||
login_params: %{
|
||||
user_id: aprs_user_id,
|
||||
passcode: aprs_passcode,
|
||||
|
|
@ -223,20 +224,29 @@ defmodule Aprs.Is do
|
|||
current_time = System.system_time(:second)
|
||||
packet_stats = update_packet_stats(state.packet_stats, current_time)
|
||||
|
||||
# Handle the incoming message
|
||||
# Task.start(Aprs, :dispatch, [packet])
|
||||
# Append new packet data to buffer
|
||||
buffer = state.buffer <> packet
|
||||
|
||||
if String.contains?(packet, "\n") or String.contains?(packet, "\r") do
|
||||
packet
|
||||
|> String.split("\r\n")
|
||||
|> Enum.each(&dispatch(String.trim(&1)))
|
||||
else
|
||||
dispatch(packet)
|
||||
end
|
||||
# Process complete lines (ending with \r\n or \n)
|
||||
{complete_lines, remaining_buffer} = extract_complete_lines(buffer)
|
||||
|
||||
# Dispatch each complete line
|
||||
Enum.each(complete_lines, fn line ->
|
||||
trimmed = String.trim(line)
|
||||
|
||||
if trimmed != "" do
|
||||
dispatch(trimmed)
|
||||
end
|
||||
end)
|
||||
|
||||
# Start a new timer
|
||||
timer = Process.send_after(self(), :aprs_no_message_timeout, @aprs_timeout)
|
||||
state = state |> Map.put(:timer, timer) |> Map.put(:packet_stats, packet_stats)
|
||||
|
||||
state =
|
||||
state
|
||||
|> Map.put(:timer, timer)
|
||||
|> Map.put(:packet_stats, packet_stats)
|
||||
|> Map.put(:buffer, remaining_buffer)
|
||||
|
||||
{:noreply, state}
|
||||
end
|
||||
|
|
@ -251,11 +261,43 @@ defmodule Aprs.Is do
|
|||
{:stop, :normal, state}
|
||||
end
|
||||
|
||||
# Extract complete lines from buffer, returning {complete_lines, remaining_buffer}
|
||||
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
|
||||
if Map.has_key?(state, :buffer) and state.buffer != "" do
|
||||
Logger.warning("Terminating with incomplete packet in buffer: #{inspect(state.buffer)}")
|
||||
end
|
||||
|
||||
# Cancel timers
|
||||
if Map.has_key?(state, :timer), do: Process.cancel_timer(state.timer)
|
||||
if Map.has_key?(state, :keepalive_timer), do: Process.cancel_timer(state.keepalive_timer)
|
||||
|
|
|
|||
|
|
@ -861,48 +861,55 @@ defmodule AprsWeb.MapLive.Index do
|
|||
|
||||
# Only include packets with valid position data
|
||||
if lat && lng do
|
||||
# Get symbol information from data_extended
|
||||
symbol_table_id =
|
||||
get_in(data_extended, ["symbol_table_id"]) ||
|
||||
packet.symbol_table_id || "/"
|
||||
build_packet_map(packet, lat, lng, data_extended)
|
||||
end
|
||||
end
|
||||
|
||||
symbol_code =
|
||||
get_in(data_extended, ["symbol_code"]) ||
|
||||
packet.symbol_code || ">"
|
||||
defp build_packet_map(packet, lat, lng, data_extended) do
|
||||
{final_table_id, final_symbol_code} = get_validated_symbol(packet, data_extended)
|
||||
callsign = generate_callsign(packet)
|
||||
|
||||
# Validate symbol
|
||||
{final_table_id, final_symbol_code} =
|
||||
if AprsSymbols.valid_symbol?(symbol_table_id, symbol_code) do
|
||||
{symbol_table_id, symbol_code}
|
||||
else
|
||||
AprsSymbols.default_symbol()
|
||||
end
|
||||
%{
|
||||
"id" => callsign,
|
||||
"callsign" => callsign,
|
||||
"base_callsign" => packet.base_callsign || "",
|
||||
"ssid" => packet.ssid || "",
|
||||
"lat" => lat,
|
||||
"lng" => lng,
|
||||
"symbol_table_id" => final_table_id,
|
||||
"symbol_code" => final_symbol_code,
|
||||
"symbol_description" => AprsSymbols.symbol_description(final_table_id, final_symbol_code),
|
||||
"data_type" => to_string(packet.data_type || "unknown"),
|
||||
"path" => packet.path || "",
|
||||
"comment" => get_in(data_extended, ["comment"]) || "",
|
||||
"data_extended" => data_extended || %{}
|
||||
}
|
||||
end
|
||||
|
||||
# Generate callsign for display
|
||||
callsign =
|
||||
if packet.ssid && packet.ssid != "" do
|
||||
"#{packet.base_callsign}-#{packet.ssid}"
|
||||
else
|
||||
packet.base_callsign || ""
|
||||
end
|
||||
defp get_validated_symbol(packet, data_extended) do
|
||||
symbol_table_id = get_symbol_table_id(packet, data_extended)
|
||||
symbol_code = get_symbol_code(packet, data_extended)
|
||||
|
||||
%{
|
||||
"id" => callsign,
|
||||
"callsign" => callsign,
|
||||
"base_callsign" => packet.base_callsign || "",
|
||||
"ssid" => packet.ssid || "",
|
||||
"lat" => lat,
|
||||
"lng" => lng,
|
||||
"symbol_table_id" => final_table_id,
|
||||
"symbol_code" => final_symbol_code,
|
||||
"symbol_description" => AprsSymbols.symbol_description(final_table_id, final_symbol_code),
|
||||
"data_type" => to_string(packet.data_type || "unknown"),
|
||||
"path" => packet.path || "",
|
||||
"comment" => get_in(data_extended, ["comment"]) || "",
|
||||
"data_extended" => data_extended || %{}
|
||||
}
|
||||
if AprsSymbols.valid_symbol?(symbol_table_id, symbol_code) do
|
||||
{symbol_table_id, symbol_code}
|
||||
else
|
||||
AprsSymbols.default_symbol()
|
||||
end
|
||||
end
|
||||
|
||||
# Return nil for packets without position data
|
||||
defp get_symbol_table_id(packet, data_extended) do
|
||||
get_in(data_extended, ["symbol_table_id"]) || packet.symbol_table_id || "/"
|
||||
end
|
||||
|
||||
defp get_symbol_code(packet, data_extended) do
|
||||
get_in(data_extended, ["symbol_code"]) || packet.symbol_code || ">"
|
||||
end
|
||||
|
||||
defp generate_callsign(packet) do
|
||||
if packet.ssid && packet.ssid != "" do
|
||||
"#{packet.base_callsign}-#{packet.ssid}"
|
||||
else
|
||||
packet.base_callsign || ""
|
||||
end
|
||||
end
|
||||
|
||||
|
|
|
|||
|
|
@ -21,7 +21,6 @@ defmodule Parser do
|
|||
|
||||
{:ok,
|
||||
%{
|
||||
# TODO: temporary for liveview
|
||||
id: 16 |> :crypto.strong_rand_bytes() |> Base.encode16(case: :lower),
|
||||
sender: sender,
|
||||
path: path,
|
||||
|
|
|
|||
|
|
@ -72,34 +72,38 @@ defmodule Parser.Types.MicE do
|
|||
Gets a value and updates it with the given function.
|
||||
"""
|
||||
def get_and_update(mic_e, key, fun) do
|
||||
value =
|
||||
case key do
|
||||
:latitude ->
|
||||
if is_number(mic_e.lat_degrees) and is_number(mic_e.lat_minutes) do
|
||||
lat = mic_e.lat_degrees + mic_e.lat_minutes / 60.0
|
||||
if mic_e.lat_direction == :south, do: -lat, else: lat
|
||||
end
|
||||
value = get_value(mic_e, key)
|
||||
apply_update_function(mic_e, key, value, fun)
|
||||
end
|
||||
|
||||
:longitude ->
|
||||
if is_number(mic_e.lon_degrees) and is_number(mic_e.lon_minutes) do
|
||||
lon = mic_e.lon_degrees + mic_e.lon_minutes / 60.0
|
||||
if mic_e.lon_direction == :west, do: -lon, else: lon
|
||||
end
|
||||
defp get_value(mic_e, :latitude), do: calculate_latitude(mic_e)
|
||||
defp get_value(mic_e, :longitude), do: calculate_longitude(mic_e)
|
||||
defp get_value(mic_e, key) when is_binary(key), do: get_string_key_value(mic_e, key)
|
||||
defp get_value(mic_e, key), do: Map.get(mic_e, key)
|
||||
|
||||
key when is_binary(key) ->
|
||||
# Handle string keys by converting to atom if it exists
|
||||
try do
|
||||
atom_key = String.to_existing_atom(key)
|
||||
Map.get(mic_e, atom_key)
|
||||
rescue
|
||||
ArgumentError ->
|
||||
nil
|
||||
end
|
||||
defp calculate_latitude(mic_e) do
|
||||
if is_number(mic_e.lat_degrees) and is_number(mic_e.lat_minutes) do
|
||||
lat = mic_e.lat_degrees + mic_e.lat_minutes / 60.0
|
||||
if mic_e.lat_direction == :south, do: -lat, else: lat
|
||||
end
|
||||
end
|
||||
|
||||
_ ->
|
||||
Map.get(mic_e, key)
|
||||
end
|
||||
defp calculate_longitude(mic_e) do
|
||||
if is_number(mic_e.lon_degrees) and is_number(mic_e.lon_minutes) do
|
||||
lon = mic_e.lon_degrees + mic_e.lon_minutes / 60.0
|
||||
if mic_e.lon_direction == :west, do: -lon, else: lon
|
||||
end
|
||||
end
|
||||
|
||||
defp get_string_key_value(mic_e, key) do
|
||||
atom_key = String.to_existing_atom(key)
|
||||
Map.get(mic_e, atom_key)
|
||||
rescue
|
||||
ArgumentError ->
|
||||
nil
|
||||
end
|
||||
|
||||
defp apply_update_function(mic_e, key, value, fun) do
|
||||
case fun.(value) do
|
||||
{get, update} -> {get, Map.put(mic_e, key, update)}
|
||||
:pop -> {value, Map.put(mic_e, key, nil)}
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue