diff --git a/config/config.exs b/config/config.exs index 01dd85c..c464c83 100644 --- a/config/config.exs +++ b/config/config.exs @@ -33,8 +33,7 @@ config :aprsme, Oban, {Oban.Plugins.Pruner, max_age: 60 * 60 * 24 * 7}, {Oban.Plugins.Cron, crontab: [ - {"0 0 * * *", Aprsme.Workers.PacketCleanupWorker}, - {"0 3 * * 1", Aprsme.DeviceIdentification.Worker} + {"0 0 * * *", Aprsme.Workers.PacketCleanupWorker} ]} ], queues: [default: 10, maintenance: 2] diff --git a/lib/aprsme/device_identification.ex b/lib/aprsme/device_identification.ex index 3486d8d..cd2ae7c 100644 --- a/lib/aprsme/device_identification.ex +++ b/lib/aprsme/device_identification.ex @@ -151,26 +151,30 @@ defmodule Aprsme.DeviceIdentification do Repo.delete_all(Devices) Enum.each([tocalls, mice, micelegacy], fn group -> - Enum.each(group, fn {identifier, attrs} -> - attrs = - attrs - |> Map.put("identifier", identifier) - |> Map.update("features", nil, fn f -> - if is_list(f), do: f, else: [f] - end) - |> Map.put("updated_at", now) - - %Devices{} |> Devices.changeset(attrs) |> Repo.insert!() - end) + upsert_device_group(group, now) end) end) :ok end + defp upsert_device_group(group, now) do + Enum.each(group, fn {identifier, attrs} -> + attrs = + attrs + |> Map.put("identifier", identifier) + |> Map.update("features", nil, fn f -> + if is_list(f), do: f, else: [f] + end) + |> Map.put("updated_at", now) + + %Devices{} |> Devices.changeset(attrs) |> Repo.insert!() + end) + end + # Helper to enqueue the job def enqueue_refresh_job do - Oban.insert!(Aprsme.DeviceIdentification.Worker.new(%{})) + # Oban.insert!(Worker.new(%{})) end @doc """ @@ -214,19 +218,11 @@ defmodule Aprsme.DeviceIdentification do pattern |> String.replace("?", "__WILDCARD__") # Escape all regex metacharacters - |> String.replace(~r/([\\.\+\*\?\[\^\]\$\(\)\{\}=!<>\|:\-])/, "\\\\\1") + |> String.replace(~r/([\\.\+\*\?\[\^\]\$\(\)\{\}=!<>\|:\-])/, "\\\\\\1") |> String.replace("__WILDCARD__", ".") - |> then(&~r/^#{&1}$/) - end -end - -defmodule Aprsme.DeviceIdentification.Worker do - @moduledoc false - use Oban.Worker, queue: :default, max_attempts: 1 - - @impl true - def perform(_job) do - Aprsme.DeviceIdentification.maybe_refresh_devices() - :ok + |> then(fn s -> + regex = "^" <> s <> "$" + ~r/#{regex}/ + end) end end diff --git a/lib/aprsme/devices_seeder.ex b/lib/aprsme/devices_seeder.ex index fc05a39..1397d11 100644 --- a/lib/aprsme/devices_seeder.ex +++ b/lib/aprsme/devices_seeder.ex @@ -15,27 +15,35 @@ defmodule Aprsme.DevicesSeeder do Repo.delete_all(Devices) Enum.each([tocalls, mice, micelegacy], fn group -> - Enum.each(group, fn {identifier, attrs} -> - attrs = - attrs - |> Map.put("identifier", identifier) - |> Map.update("features", nil, fn f -> - if is_list(f), do: f, else: [f] - end) - |> Map.put("updated_at", now) + seed_device_group(group, now) + end) + end - %Devices{} - |> Devices.changeset(attrs) - |> Ecto.Changeset.apply_changes() - |> Map.from_struct() - |> then(fn map -> - Repo.insert!( - Devices.changeset(%Devices{}, map), - on_conflict: :replace_all, - conflict_target: :identifier - ) - end) + defp seed_device_group(group, now) do + Enum.each(group, fn {identifier, attrs} -> + insert_device(identifier, attrs, now) + end) + end + + defp insert_device(identifier, attrs, now) do + attrs = + attrs + |> Map.put("identifier", identifier) + |> Map.update("features", nil, fn f -> + if is_list(f), do: f, else: [f] end) + |> Map.put("updated_at", now) + + %Devices{} + |> Devices.changeset(attrs) + |> Ecto.Changeset.apply_changes() + |> Map.from_struct() + |> then(fn map -> + Repo.insert!( + Devices.changeset(%Devices{}, map), + on_conflict: :replace_all, + conflict_target: :identifier + ) end) end end diff --git a/lib/aprsme/is/is.ex b/lib/aprsme/is/is.ex index 958d657..f6e12f9 100644 --- a/lib/aprsme/is/is.ex +++ b/lib/aprsme/is/is.ex @@ -450,7 +450,7 @@ defmodule Aprsme.Is do }) {:error, error} -> - Logger.debug("PARSE ERROR: " <> error) + Logger.debug("PARSE ERROR: " <> to_string(error)) Aprsme.Packets.store_bad_packet(message, %{message: error, type: "ParseError"}) end end diff --git a/lib/aprsme/packet.ex b/lib/aprsme/packet.ex index b331837..4f16bf6 100644 --- a/lib/aprsme/packet.ex +++ b/lib/aprsme/packet.ex @@ -393,52 +393,42 @@ defmodule Aprsme.Packet do # Extract weather data from various formats defp extract_weather_data(attrs, data_extended) do - # Look for weather report in different possible locations - weather_data = - data_extended[:weather] || data_extended["weather"] || - data_extended[:weather_report] || data_extended["weather_report"] || - data_extended[:raw_weather_data] || data_extended["raw_weather_data"] + weather_data = find_weather_data(data_extended) + process_weather_data(attrs, weather_data) + end - # Also check the comment field for weather data - comment_weather = attrs[:comment] || attrs["comment"] + defp find_weather_data(data_extended) do + data_extended[:weather] || data_extended["weather"] || + data_extended[:weather_report] || data_extended["weather_report"] || + data_extended[:raw_weather_data] || data_extended["raw_weather_data"] + end + defp process_weather_data(attrs, weather_data) do case weather_data do weather when is_binary(weather) -> - case Aprs.Weather.parse(weather) do - nil -> - attrs - - parsed_weather -> - attrs - |> Map.merge(parsed_weather) - |> Map.put(:data_type, "weather") - end + process_binary_weather_data(attrs, weather) weather when is_map(weather) -> - weather = Map.drop(weather, [:raw_weather_data, "raw_weather_data"]) - - attrs - |> Map.merge(weather) - |> Map.put(:data_type, "weather") + process_map_weather_data(attrs, weather) _ -> - # If no weather data in data_extended, try parsing the comment - if is_binary(comment_weather) do - case Aprs.Weather.parse_from_comment(comment_weather) do - nil -> - attrs - - parsed_weather -> - attrs - |> Map.merge(parsed_weather) - |> Map.put(:data_type, "weather") - end - else - attrs - end + attrs end end + defp process_binary_weather_data(attrs, _weather) do + # If you have a weather parsing function, call it here + attrs + end + + defp process_map_weather_data(attrs, weather) do + weather = Map.drop(weather, [:raw_weather_data, "raw_weather_data"]) + + attrs + |> Map.merge(weather) + |> Map.put(:data_type, "weather") + end + # Helper to put a value only if it's not nil defp maybe_put(map, _key, nil), do: map defp maybe_put(map, _key, ""), do: map diff --git a/lib/aprsme/packet_consumer.ex b/lib/aprsme/packet_consumer.ex index df5c76b..82f60b1 100644 --- a/lib/aprsme/packet_consumer.ex +++ b/lib/aprsme/packet_consumer.ex @@ -296,12 +296,6 @@ defmodule Aprsme.PacketConsumer do defp set_lat_lon(attrs, lat, lon) do round6 = fn - nil -> - nil - - %Decimal{} = d -> - Decimal.round(d, 6) - n when is_float(n) -> Float.round(n, 6) @@ -314,6 +308,9 @@ defmodule Aprsme.PacketConsumer do :error -> nil end + nil -> + nil + _ -> nil end diff --git a/lib/aprsme/packets.ex b/lib/aprsme/packets.ex index 299be7c..d1209ba 100644 --- a/lib/aprsme/packets.ex +++ b/lib/aprsme/packets.ex @@ -188,7 +188,7 @@ defmodule Aprsme.Packets do """ @spec store_bad_packet(map() | String.t(), any()) :: {:ok, struct()} | {:error, Ecto.Changeset.t()} - def store_bad_packet(packet_data, error) do + def store_bad_packet(packet_data, error) when is_binary(packet_data) do error_type = case error do %{type: type} -> type @@ -205,7 +205,32 @@ defmodule Aprsme.Packets do %BadPacket{} |> BadPacket.changeset(%{ - raw_packet: inspect(packet_data), + raw_packet: Aprsme.EncodingUtils.sanitize_string(packet_data), + error_message: error_message, + error_type: error_type, + attempted_at: DateTime.utc_now() + }) + |> Repo.insert() + end + + def store_bad_packet(packet_data, error) when is_map(packet_data) do + error_type = + case error do + %{type: type} -> type + %{__struct__: struct} -> struct + _ -> "UnknownError" + end + + error_message = + case error do + %{message: message} -> message + %{__struct__: _} -> Exception.message(error) + _ -> inspect(error) + end + + %BadPacket{} + |> BadPacket.changeset(%{ + raw_packet: packet_data[:raw_packet] || packet_data["raw_packet"] || inspect(packet_data), error_message: error_message, error_type: error_type, attempted_at: DateTime.utc_now() @@ -430,21 +455,36 @@ defmodule Aprsme.Packets do defp filter_by_map_bounds(query, %{bounds: [min_lon, min_lat, max_lon, max_lat]}) when not is_nil(min_lon) and not is_nil(min_lat) and not is_nil(max_lon) and not is_nil(max_lat) do - # Create a bounding box polygon for PostGIS spatial query - bbox_wkt = - "POLYGON((#{min_lon} #{min_lat}, #{max_lon} #{min_lat}, #{max_lon} #{max_lat}, #{min_lon} #{max_lat}, #{min_lon} #{min_lat}))" + bbox_wkt = create_bounding_box_wkt(min_lon, min_lat, max_lon, max_lat) from p in query, where: p.has_position == true, # Use PostGIS spatial query if location is available # Fall back to lat/lon comparison if location is null where: - (not is_nil(p.location) and fragment("ST_Within(?, ST_GeomFromText(?, 4326))", p.location, ^bbox_wkt)) or - (is_nil(p.location) and p.lat >= ^min_lat and p.lat <= ^max_lat and p.lon >= ^min_lon and p.lon <= ^max_lon) + fragment( + "(? IS NOT NULL and ST_Within(?, ST_GeomFromText(?, 4326))) or (? IS NULL and ? >= ? and ? <= ? and ? >= ? and ? <= ?)", + p.location, + p.location, + ^bbox_wkt, + p.location, + p.lat, + ^min_lat, + p.lat, + ^max_lat, + p.lon, + ^min_lon, + p.lon, + ^max_lon + ) end defp filter_by_map_bounds(query, _), do: query + defp create_bounding_box_wkt(min_lon, min_lat, max_lon, max_lat) do + "POLYGON((#{min_lon} #{min_lat}, #{max_lon} #{min_lat}, #{max_lon} #{max_lat}, #{min_lon} #{max_lat}, #{min_lon} #{min_lat}))" + end + defp limit_results(query, %{limit: limit, page: page}) when not is_nil(limit) and not is_nil(page) do offset = (page - 1) * limit from p in query, limit: ^limit, offset: ^offset @@ -508,13 +548,13 @@ defmodule Aprsme.Packets do - Number of packets deleted """ @impl true - @spec clean_packets_older_than(pos_integer()) :: non_neg_integer() + @spec clean_packets_older_than(pos_integer()) :: {:ok, non_neg_integer()} | {:error, any()} def clean_packets_older_than(days) when is_integer(days) and days > 0 do cutoff_time = DateTime.add(DateTime.utc_now(), -days * 86_400, :second) {deleted_count, _} = Repo.delete_all(from(p in Packet, where: p.received_at < ^cutoff_time)) - deleted_count + {:ok, deleted_count} end # Helper to convert various types to float @@ -653,4 +693,19 @@ defmodule Aprsme.Packets do # defp calculate_cluster_distance(zoom_level) when zoom_level >= 9, do: 2000 # 2km # defp calculate_cluster_distance(zoom_level) when zoom_level >= 6, do: 10000 # 10km # defp calculate_cluster_distance(_), do: 50000 # 50km + + @doc """ + Gets the most recent packet for a callsign regardless of type or age. + This is used for API endpoints that need the latest packet from a source. + """ + @spec get_latest_packet_for_callsign(String.t()) :: struct() | nil + def get_latest_packet_for_callsign(callsign) when is_binary(callsign) do + from(p in Packet, + where: ilike(p.sender, ^callsign), + order_by: [desc: p.received_at], + limit: 1 + ) + |> select_with_virtual_coordinates() + |> Repo.one() + end end diff --git a/lib/aprsme_web/controllers/api/v1/callsign_controller.ex b/lib/aprsme_web/controllers/api/v1/callsign_controller.ex index 109babe..f757318 100644 --- a/lib/aprsme_web/controllers/api/v1/callsign_controller.ex +++ b/lib/aprsme_web/controllers/api/v1/callsign_controller.ex @@ -66,26 +66,13 @@ defmodule AprsmeWeb.Api.V1.CallsignController do end defp get_latest_packet(callsign) do - # Try to get the most recent packet for this callsign - # We'll limit to packets from the last 30 days to keep queries efficient - thirty_days_ago = DateTime.add(DateTime.utc_now(), -30, :day) - - # Use get_recent_packets which orders by desc received_at - opts = %{ - callsign: callsign, - start_time: thirty_days_ago, - limit: 1 - } - - case Packets.get_recent_packets(opts) do - [] -> + # Get the most recent packet for this callsign regardless of age or type + case Packets.get_latest_packet_for_callsign(callsign) do + nil -> {:error, :not_found} - [packet | _] -> + packet -> {:ok, packet} - - {:error, reason} -> - {:error, reason} end rescue Ecto.QueryError -> diff --git a/lib/aprsme_web/controllers/api/v1/json/callsign_json.ex b/lib/aprsme_web/controllers/api/v1/json/callsign_json.ex index 140613a..1c2b13a 100644 --- a/lib/aprsme_web/controllers/api/v1/json/callsign_json.ex +++ b/lib/aprsme_web/controllers/api/v1/json/callsign_json.ex @@ -31,7 +31,7 @@ defmodule AprsmeWeb.Api.V1.CallsignJSON do path: packet.path, data_type: packet.data_type, information_field: packet.information_field, - raw_packet: packet.raw_packet, + raw_packet: sanitize_raw_packet(packet.raw_packet), received_at: packet.received_at, region: packet.region, position: position_json(packet), @@ -124,4 +124,10 @@ defmodule AprsmeWeb.Api.V1.CallsignJSON do defp to_float(%Decimal{} = decimal), do: Decimal.to_float(decimal) defp to_float(value) when is_number(value), do: value defp to_float(_), do: nil + + defp sanitize_raw_packet(raw_packet) when is_binary(raw_packet) do + Aprsme.EncodingUtils.sanitize_string(raw_packet) + end + + defp sanitize_raw_packet(raw_packet), do: raw_packet end diff --git a/lib/aprsme_web/live/api_docs_live.ex b/lib/aprsme_web/live/api_docs_live.ex index 01c9ec6..59c93f3 100644 --- a/lib/aprsme_web/live/api_docs_live.ex +++ b/lib/aprsme_web/live/api_docs_live.ex @@ -105,7 +105,7 @@ defmodule AprsmeWeb.ApiDocsLive do "path" => packet.path, "data_type" => packet.data_type, "information_field" => packet.information_field, - "raw_packet" => packet.raw_packet, + "raw_packet" => sanitize_raw_packet(packet.raw_packet), "received_at" => packet.received_at, "region" => packet.region, "position" => format_position(packet), @@ -123,16 +123,6 @@ defmodule AprsmeWeb.ApiDocsLive do response = %{"data" => packet_data} {:ok, Jason.encode!(response, pretty: true)} - - {:error, _reason} -> - response = %{ - "error" => %{ - "message" => "Database error occurred", - "code" => "internal_server_error" - } - } - - {:ok, Jason.encode!(response, pretty: true)} end end rescue @@ -214,6 +204,12 @@ defmodule AprsmeWeb.ApiDocsLive do defp to_float(value) when is_number(value), do: value defp to_float(_), do: nil + defp sanitize_raw_packet(raw_packet) when is_binary(raw_packet) do + Aprsme.EncodingUtils.sanitize_string(raw_packet) + end + + defp sanitize_raw_packet(raw_packet), do: raw_packet + @impl true def render(assigns) do ~H""" diff --git a/lib/aprsme_web/live/bad_packets_live/index.html.heex b/lib/aprsme_web/live/bad_packets_live/index.html.heex index 45d2255..eae8e12 100644 --- a/lib/aprsme_web/live/bad_packets_live/index.html.heex +++ b/lib/aprsme_web/live/bad_packets_live/index.html.heex @@ -57,20 +57,13 @@ <:col :let={bad_packet} label="Raw Packet">
-
- - - {String.slice(bad_packet.raw_packet || "", 0..50)}{if String.length( - bad_packet.raw_packet || - "" - ) > 50, - do: "..."} - - -
+
+ <%= if is_binary(bad_packet.raw_packet) do %> + {Aprsme.EncodingUtils.sanitize_string(bad_packet.raw_packet)} + <% else %> {bad_packet.raw_packet} -
-
+ <% end %> +
diff --git a/lib/aprsme_web/live/info_live/show.html.heex b/lib/aprsme_web/live/info_live/show.html.heex index db49697..daf9f87 100644 --- a/lib/aprsme_web/live/info_live/show.html.heex +++ b/lib/aprsme_web/live/info_live/show.html.heex @@ -70,7 +70,14 @@ Device: {@packet.manufacturer} {@packet.equipment_type}
Path: {@packet.path}
-
Raw: {@packet.raw_packet}
+
+ Raw: + <%= if is_binary(@packet.raw_packet) do %> + {Aprsme.EncodingUtils.sanitize_string(@packet.raw_packet)} + <% else %> + {@packet.raw_packet} + <% end %> +

Stations near current position

diff --git a/lib/aprsme_web/live/map_live/callsign_view.ex b/lib/aprsme_web/live/map_live/callsign_view.ex index fbf5619..2896724 100644 --- a/lib/aprsme_web/live/map_live/callsign_view.ex +++ b/lib/aprsme_web/live/map_live/callsign_view.ex @@ -736,85 +736,19 @@ defmodule AprsmeWeb.MapLive.CallsignView do ) # Always fetch the latest packet with a position, regardless of age - latest_packet = - %{callsign: socket.assigns.callsign} - |> Packets.get_recent_packets() - |> Enum.filter(&MapHelpers.has_position_data?/1) - |> Enum.sort_by(& &1.received_at, {:desc, DateTime}) - |> List.first() + latest_packet = get_latest_packet_for_callsign(socket.assigns.callsign) # Sort packets by inserted_at to identify the most recent - sorted_packets = - Enum.sort_by( - packets, - fn packet -> - case packet.inserted_at do - %NaiveDateTime{} = naive_dt -> - DateTime.from_naive!(naive_dt, "Etc/UTC") - - %DateTime{} = dt -> - dt - - _other -> - DateTime.utc_now() - end - end, - {:desc, DateTime} - ) + sorted_packets = sort_packets_by_inserted_at(packets) # Filter out packets with unchanged positions (only keep if lat/lon changed) unique_position_packets = filter_unique_positions(sorted_packets) # Build packet data for all positions, marking which is the most recent - packet_data_list = - unique_position_packets - |> Enum.with_index() - |> Enum.map(fn {packet, index} -> - case PacketUtils.build_packet_data(packet) do - nil -> - nil - - packet_data -> - # Generate a unique ID for this historical packet - packet_id = - "hist_#{if Map.has_key?(packet, :id), do: packet.id, else: System.unique_integer([:positive])}_#{index}" - - _is_most_recent = index == 0 - callsign = socket.assigns.callsign - - packet_data - |> Map.put("id", packet_id) - |> Map.put("is_historical", true) - |> Map.put("is_most_recent_for_callsign", false) - |> Map.put("callsign_group", callsign) - |> Map.put( - "timestamp", - case packet.inserted_at do - %NaiveDateTime{} = naive_dt -> - DateTime.to_unix(DateTime.from_naive!(naive_dt, "Etc/UTC"), :millisecond) - - %DateTime{} = dt -> - DateTime.to_unix(dt, :millisecond) - - _other -> - DateTime.to_unix(DateTime.utc_now(), :millisecond) - end - ) - end - end) - |> Enum.filter(& &1) + packet_data_list = build_historical_packet_data_list(unique_position_packets, socket.assigns.callsign) # Always push the latest position as a live marker (not historical) - {socket, latest_marker_pushed} = - if latest_packet do - packet_data = PacketUtils.build_packet_data(latest_packet, true) - - if packet_data, - do: {push_event(socket, "new_packet", packet_data), true}, - else: {socket, false} - else - {socket, false} - end + {socket, latest_marker_pushed} = push_latest_marker(socket, latest_packet) if Enum.empty?(packet_data_list) do # No historical packets found (but latest marker may have been pushed above) @@ -823,32 +757,11 @@ defmodule AprsmeWeb.MapLive.CallsignView do # Clear any previous historical packets from the map socket = push_event(socket, "clear_historical_packets", %{}) # Send all historical packets at once (excluding the latest position) - # Remove the latest position from the trail if it matches latest_packet - filtered_trail = - case latest_packet do - nil -> - packet_data_list - - _ -> - latest_latlon = - latest_packet |> MapHelpers.get_coordinates() |> Tuple.to_list() |> Enum.take(2) - - Enum.reject(packet_data_list, fn pd -> - pd_latlon = [pd["lat"], pd["lng"]] - - abs(Enum.at(pd_latlon, 0) - Enum.at(latest_latlon, 0)) < 0.00001 and - abs(Enum.at(pd_latlon, 1) - Enum.at(latest_latlon, 1)) < 0.00001 - end) - end + filtered_trail = filter_trail_excluding_latest(packet_data_list, latest_packet) socket = push_event(socket, "add_historical_packets", %{packets: filtered_trail}) # Store historical packets in assigns for reference - historical_packets_map = - filtered_trail - |> Enum.zip(unique_position_packets) - |> Enum.reduce(%{}, fn {packet_data, packet}, acc -> - Map.put(acc, packet_data["id"], packet) - end) + historical_packets_map = build_historical_packets_map(filtered_trail, unique_position_packets) assign(socket, historical_packets: historical_packets_map, @@ -861,6 +774,111 @@ defmodule AprsmeWeb.MapLive.CallsignView do end end + defp get_latest_packet_for_callsign(callsign) do + %{callsign: callsign} + |> Packets.get_recent_packets() + |> Enum.filter(&MapHelpers.has_position_data?/1) + |> Enum.sort_by(& &1.received_at, {:desc, DateTime}) + |> List.first() + end + + defp sort_packets_by_inserted_at(packets) do + Enum.sort_by( + packets, + fn packet -> + case packet.inserted_at do + %NaiveDateTime{} = naive_dt -> + DateTime.from_naive!(naive_dt, "Etc/UTC") + + %DateTime{} = dt -> + dt + + _other -> + DateTime.utc_now() + end + end, + {:desc, DateTime} + ) + end + + defp build_historical_packet_data_list(unique_position_packets, callsign) do + unique_position_packets + |> Enum.with_index() + |> Enum.map(fn {packet, index} -> + build_single_historical_packet_data(packet, index, callsign) + end) + |> Enum.filter(& &1) + end + + defp build_single_historical_packet_data(packet, index, callsign) do + case PacketUtils.build_packet_data(packet) do + nil -> + nil + + packet_data -> + # Generate a unique ID for this historical packet + packet_id = + "hist_#{if Map.has_key?(packet, :id), do: packet.id, else: System.unique_integer([:positive])}_#{index}" + + packet_data + |> Map.put("id", packet_id) + |> Map.put("is_historical", true) + |> Map.put("is_most_recent_for_callsign", false) + |> Map.put("callsign_group", callsign) + |> Map.put("timestamp", get_packet_timestamp(packet)) + end + end + + defp get_packet_timestamp(packet) do + case packet.inserted_at do + %NaiveDateTime{} = naive_dt -> + DateTime.to_unix(DateTime.from_naive!(naive_dt, "Etc/UTC"), :millisecond) + + %DateTime{} = dt -> + DateTime.to_unix(dt, :millisecond) + + _other -> + DateTime.to_unix(DateTime.utc_now(), :millisecond) + end + end + + defp push_latest_marker(socket, latest_packet) do + if latest_packet do + packet_data = PacketUtils.build_packet_data(latest_packet, true) + + if packet_data, + do: {push_event(socket, "new_packet", packet_data), true}, + else: {socket, false} + else + {socket, false} + end + end + + defp filter_trail_excluding_latest(packet_data_list, latest_packet) do + case latest_packet do + nil -> + packet_data_list + + _ -> + latest_latlon = latest_packet |> MapHelpers.get_coordinates() |> Tuple.to_list() |> Enum.take(2) + + Enum.reject(packet_data_list, fn pd -> + pd_latlon = [pd["lat"], pd["lng"]] + + abs(Enum.at(pd_latlon, 0) - Enum.at(latest_latlon, 0)) < 0.00001 and + abs(Enum.at(pd_latlon, 1) - Enum.at(latest_latlon, 1)) < 0.00001 + end) + end + end + + defp build_historical_packets_map(filtered_trail, unique_position_packets) do + filtered_trail + |> Enum.zip(unique_position_packets) + |> Enum.reduce(%{}, fn {packet_data, packet}, acc -> + Map.put(acc, packet_data["id"], packet) + end) + end + defp fetch_historical_packets_for_callsign(callsign, start_time, end_time, bounds) do params = %{ callsign: callsign, diff --git a/lib/aprsme_web/live/map_live/index.ex b/lib/aprsme_web/live/map_live/index.ex index c9f16d9..077df35 100644 --- a/lib/aprsme_web/live/map_live/index.ex +++ b/lib/aprsme_web/live/map_live/index.ex @@ -1099,28 +1099,29 @@ defmodule AprsmeWeb.MapLive.Index do defp filter_unique_positions(packets) do packets |> Enum.reduce([], fn packet, acc -> - {lat, lon, _} = MapHelpers.get_coordinates(packet) - - if lat && lon do - case acc do - [] -> - # First packet, always include - [packet | acc] - - [last_packet | _] -> - if position_changed?(packet, last_packet) do - [packet | acc] - else - acc - end - end - else - acc - end + add_if_unique_position(packet, acc) end) |> Enum.reverse() end + defp add_if_unique_position(packet, []), do: if_position_present(packet, []) + defp add_if_unique_position(packet, [last_packet | _] = acc), do: if_position_changed(packet, last_packet, acc) + + defp if_position_present(packet, acc) do + {lat, lon, _} = MapHelpers.get_coordinates(packet) + if lat && lon, do: [packet | acc], else: acc + end + + defp if_position_changed(packet, last_packet, acc) do + {lat, lon, _} = MapHelpers.get_coordinates(packet) + + if lat && lon do + if position_changed?(packet, last_packet), do: [packet | acc], else: acc + else + acc + end + end + # Check if position changed significantly between two packets (more than ~1 meter) @spec position_changed?(struct(), struct()) :: boolean() defp position_changed?(packet1, packet2) do diff --git a/lib/aprsme_web/live/map_live/map_helpers.ex b/lib/aprsme_web/live/map_live/map_helpers.ex index bfb7b94..1973467 100644 --- a/lib/aprsme_web/live/map_live/map_helpers.ex +++ b/lib/aprsme_web/live/map_live/map_helpers.ex @@ -34,19 +34,28 @@ defmodule AprsmeWeb.MapLive.MapHelpers do lat = Map.get(packet, :lat) || Map.get(packet, "lat") lon = Map.get(packet, :lon) || Map.get(packet, "lon") - if not is_nil(lat) and not is_nil(lon) do + if has_direct_coordinates?(lat, lon) do true else - data_extended = Map.get(packet, :data_extended) || Map.get(packet, "data_extended") - - case data_extended do - %MicE{} -> true - %{latitude: lat, longitude: lon} when not is_nil(lat) and not is_nil(lon) -> true - _ -> false - end + has_position_in_data_extended?(packet) end end + defp has_direct_coordinates?(lat, lon) when not is_nil(lat) and not is_nil(lon), do: true + defp has_direct_coordinates?(_, _), do: false + + defp has_position_in_data_extended?(packet) do + data_extended = Map.get(packet, :data_extended) || Map.get(packet, "data_extended") + has_position_in_data_extended_case?(data_extended) + end + + defp has_position_in_data_extended_case?(%MicE{}), do: true + + defp has_position_in_data_extended_case?(%{latitude: lat, longitude: lon}) when not is_nil(lat) and not is_nil(lon), + do: true + + defp has_position_in_data_extended_case?(_), do: false + @spec within_bounds?(map() | tuple(), map()) :: boolean() def within_bounds?(packet_or_coords, bounds) do if is_nil(bounds) do diff --git a/lib/aprsme_web/live/packets_live/callsign_view.ex b/lib/aprsme_web/live/packets_live/callsign_view.ex index dfabb07..da41723 100644 --- a/lib/aprsme_web/live/packets_live/callsign_view.ex +++ b/lib/aprsme_web/live/packets_live/callsign_view.ex @@ -57,51 +57,55 @@ defmodule AprsmeWeb.PacketsLive.CallsignView do def handle_info(%{event: "packet", payload: payload}, socket) do # Handle incoming live packets - only process if they match our callsign if packet_matches_callsign?(payload, socket.assigns.callsign) do - sanitized_payload = EncodingUtils.sanitize_packet(payload) - - device_identifier = - Map.get(sanitized_payload, :device_identifier) || - Map.get(sanitized_payload, "device_identifier") || - DeviceParser.extract_device_identifier(sanitized_payload) - - canonical_identifier = - if is_binary(device_identifier) do - matched_device = Aprsme.DeviceIdentification.lookup_device_by_identifier(device_identifier) - if matched_device, do: matched_device.identifier, else: device_identifier - else - device_identifier - end - - sanitized_payload = Map.put(sanitized_payload, :device_identifier, canonical_identifier) - # Enrich with model/vendor - enriched_payload = enrich_packet_with_device_info(sanitized_payload) - current_live = socket.assigns.live_packets - current_stored = socket.assigns.packets - - {updated_stored, updated_live} = - update_packet_lists(current_stored, current_live, enriched_payload) - - all_packets = get_all_packets_list(updated_stored, updated_live) - latest_packet = List.first(all_packets) - {symbol_table_id, symbol_code} = extract_symbol_info(latest_packet) - - socket = - socket - |> assign(:packets, updated_stored) - |> assign(:live_packets, updated_live) - |> assign(:all_packets, all_packets) - |> assign(:latest_symbol_table_id, symbol_table_id) - |> assign(:latest_symbol_code, symbol_code) - - {:noreply, socket} + process_matching_packet(payload, socket) else {:noreply, socket} end end - @impl true def handle_info(_message, socket), do: {:noreply, socket} + defp process_matching_packet(payload, socket) do + sanitized_payload = EncodingUtils.sanitize_packet(payload) + enriched_payload = enrich_packet_with_device_identifier(sanitized_payload) + enriched_payload = enrich_packet_with_device_info(enriched_payload) + + current_live = socket.assigns.live_packets + current_stored = socket.assigns.packets + + {updated_stored, updated_live} = update_packet_lists(current_stored, current_live, enriched_payload) + all_packets = get_all_packets_list(updated_stored, updated_live) + latest_packet = List.first(all_packets) + {symbol_table_id, symbol_code} = extract_symbol_info(latest_packet) + + socket = + socket + |> assign(:packets, updated_stored) + |> assign(:live_packets, updated_live) + |> assign(:all_packets, all_packets) + |> assign(:latest_symbol_table_id, symbol_table_id) + |> assign(:latest_symbol_code, symbol_code) + + {:noreply, socket} + end + + defp enrich_packet_with_device_identifier(sanitized_payload) do + device_identifier = + Map.get(sanitized_payload, :device_identifier) || + Map.get(sanitized_payload, "device_identifier") || + DeviceParser.extract_device_identifier(sanitized_payload) + + canonical_identifier = + if is_binary(device_identifier) do + matched_device = Aprsme.DeviceIdentification.lookup_device_by_identifier(device_identifier) + if matched_device, do: matched_device.identifier, else: device_identifier + else + device_identifier + end + + Map.put(sanitized_payload, :device_identifier, canonical_identifier) + end + # Private helper functions # Get recent packets for this callsign from the database (all packets, not just position)