Turned on :error_handling, :underspecs, and :unmatched_returns in mix.exs dialyzer config. The 97 warnings this surfaced were fixed in place rather than suppressed: - unmatched_return (79): explicit discard with `_ = ...` for fire-and-forget side effects (Process.cancel_timer, :ets.new, send/2), and pattern-matched `:ok = ...` for control-plane Phoenix.PubSub subscribe/unsubscribe/broadcast calls so a future return-shape change fails loud. - contract_supertype (18): tightened @spec arg and return types on data_builder, historical_loader, url_params, packet_utils, encoding_utils, aprs_symbol, weather_controller, packet_replay to match each function's actual success typing. No behavioural change. mix compile clean, 1008 tests pass, dialyzer count is now 0.
776 lines
26 KiB
Elixir
776 lines
26 KiB
Elixir
defmodule Aprsme.PacketConsumer do
|
|
@moduledoc """
|
|
GenStage consumer that batches APRS packets and inserts them into the database
|
|
efficiently to reduce database load.
|
|
"""
|
|
use GenStage
|
|
|
|
alias Aprsme.Cluster.PacketDistributor
|
|
alias Aprsme.LogSanitizer
|
|
alias Aprsme.Repo
|
|
|
|
require Logger
|
|
|
|
@type state :: %{
|
|
batch: list(map()),
|
|
batch_length: non_neg_integer(),
|
|
batch_size: integer(),
|
|
batch_timeout: integer(),
|
|
max_batch_size: integer(),
|
|
timer: reference() | nil
|
|
}
|
|
|
|
@spec start_link(keyword()) :: GenServer.on_start()
|
|
def start_link(opts \\ []) do
|
|
# Allow unnamed consumers for pool usage
|
|
name = opts[:name]
|
|
|
|
if name do
|
|
GenStage.start_link(__MODULE__, opts, name: name)
|
|
else
|
|
GenStage.start_link(__MODULE__, opts)
|
|
end
|
|
end
|
|
|
|
@impl true
|
|
def init(opts) do
|
|
batch_size = opts[:batch_size] || 100
|
|
batch_timeout = opts[:batch_timeout] || 1000
|
|
# Maximum batch size to prevent unbounded memory growth
|
|
max_batch_size = opts[:max_batch_size] || 1000
|
|
|
|
# Start a timer for batch processing
|
|
timer = Process.send_after(self(), :process_batch, batch_timeout)
|
|
|
|
# Extract subscription options if provided
|
|
subscribe_to = opts[:subscribe_to] || [{Aprsme.PacketProducer, max_demand: opts[:max_demand] || 250}]
|
|
|
|
{:consumer,
|
|
%{
|
|
batch: [],
|
|
batch_length: 0,
|
|
batch_size: batch_size,
|
|
batch_timeout: batch_timeout,
|
|
max_batch_size: max_batch_size,
|
|
timer: timer
|
|
}, subscribe_to: subscribe_to}
|
|
end
|
|
|
|
@impl true
|
|
def handle_events(events, _from, state) do
|
|
handle_batch_update(events, state)
|
|
end
|
|
|
|
# Pattern matching for batch handling with optimized list operations
|
|
defp handle_batch_update(events, %{batch: batch, batch_length: batch_length} = state) do
|
|
# Optimize: Keep batch in reverse order for O(1) prepending
|
|
# Only reverse when processing
|
|
new_batch = Enum.reverse(events, batch)
|
|
new_batch_length = batch_length + length(events)
|
|
|
|
handle_batch_by_size(new_batch, new_batch_length, state)
|
|
end
|
|
|
|
# Pattern matching for different batch size scenarios
|
|
defp handle_batch_by_size(batch, length, %{max_batch_size: max} = state) when length >= max do
|
|
handle_oversized_batch(batch, max, state)
|
|
end
|
|
|
|
defp handle_batch_by_size(batch, length, %{batch_size: size} = state) when length >= size do
|
|
handle_full_batch(batch, state)
|
|
end
|
|
|
|
defp handle_batch_by_size(batch, new_batch_length, state) do
|
|
handle_partial_batch(batch, new_batch_length, state)
|
|
end
|
|
|
|
# Handle oversized batch with pattern matching.
|
|
# Process the oldest `max_size` packets now and keep the rest for the next
|
|
# batch cycle — dropping them here would silently lose packets that the
|
|
# producer has already acknowledged.
|
|
defp handle_oversized_batch(batch, max_size, state) do
|
|
# Reverse once to process in arrival order
|
|
reversed_batch = Enum.reverse(batch)
|
|
{process_now, remainder} = Enum.split(reversed_batch, max_size)
|
|
|
|
process_batch(process_now)
|
|
|
|
# Keep the remainder queued. Internal batch is stored reversed
|
|
# (newest-at-head) so re-reverse the leftover arrival-order slice.
|
|
carryover_batch = Enum.reverse(remainder)
|
|
carryover_length = length(remainder)
|
|
|
|
if carryover_length > 0 do
|
|
Logger.info("Carrying over #{carryover_length} packets past batch-size limit to next cycle",
|
|
batch_info:
|
|
LogSanitizer.log_data(
|
|
carryover_count: carryover_length,
|
|
processed_count: length(process_now)
|
|
)
|
|
)
|
|
end
|
|
|
|
new_state =
|
|
state
|
|
|> reset_batch_timer()
|
|
|> Map.put(:batch, carryover_batch)
|
|
|> Map.put(:batch_length, carryover_length)
|
|
|
|
{:noreply, [], new_state}
|
|
end
|
|
|
|
# Handle full batch
|
|
defp handle_full_batch(batch, state) do
|
|
# Reverse once for processing
|
|
process_batch(Enum.reverse(batch))
|
|
|
|
new_state = reset_batch_timer(state)
|
|
{:noreply, [], new_state}
|
|
end
|
|
|
|
# Handle partial batch
|
|
defp handle_partial_batch(batch, new_batch_length, state) do
|
|
# Keep batch in reverse order
|
|
{:noreply, [], %{state | batch: batch, batch_length: new_batch_length}}
|
|
end
|
|
|
|
# Helper functions with pattern matching
|
|
@spec reset_batch_timer(state()) :: state()
|
|
defp reset_batch_timer(%{timer: nil} = state) do
|
|
new_timer = Process.send_after(self(), :process_batch, state.batch_timeout)
|
|
%{state | batch: [], batch_length: 0, timer: new_timer}
|
|
end
|
|
|
|
defp reset_batch_timer(%{timer: timer} = state) do
|
|
_ = Process.cancel_timer(timer)
|
|
new_timer = Process.send_after(self(), :process_batch, state.batch_timeout)
|
|
%{state | batch: [], batch_length: 0, timer: new_timer}
|
|
end
|
|
|
|
@impl true
|
|
# Pattern matching for empty batch
|
|
def handle_info(:process_batch, %{batch: []} = state) do
|
|
# Just restart the timer for empty batch
|
|
new_state = start_batch_timer(state)
|
|
{:noreply, [], new_state}
|
|
end
|
|
|
|
# Pattern matching for non-empty batch
|
|
def handle_info(:process_batch, %{batch: batch, batch_length: batch_length} = state) when is_list(batch) do
|
|
check_batch_utilization(batch_length, state.max_batch_size)
|
|
# Remember to reverse the batch since we keep it in reverse order
|
|
process_batch(Enum.reverse(batch))
|
|
|
|
new_state = start_batch_timer(%{state | batch: [], batch_length: 0})
|
|
{:noreply, [], new_state}
|
|
end
|
|
|
|
# Helper to start batch timer
|
|
@spec start_batch_timer(state()) :: state()
|
|
defp start_batch_timer(%{batch_timeout: timeout} = state) do
|
|
timer = Process.send_after(self(), :process_batch, timeout)
|
|
%{state | timer: timer}
|
|
end
|
|
|
|
# Batch utilization check using pre-computed length to avoid O(n) traversals
|
|
@spec check_batch_utilization(non_neg_integer(), integer()) :: :ok
|
|
defp check_batch_utilization(batch_length, max_size) when batch_length > max_size * 0.8 do
|
|
Logger.warning("Batch size approaching limit",
|
|
batch_status:
|
|
LogSanitizer.log_data(
|
|
current_size: batch_length,
|
|
max_size: max_size,
|
|
utilization_percent: trunc(batch_length / max_size * 100)
|
|
)
|
|
)
|
|
end
|
|
|
|
defp check_batch_utilization(_batch_length, _max_size), do: :ok
|
|
|
|
@spec process_batch(list(map())) :: :ok
|
|
defp process_batch(packets) do
|
|
require Logger
|
|
|
|
# Monitor memory usage before processing (only this process)
|
|
{:memory, memory_before} = Process.info(self(), :memory)
|
|
start_time = System.monotonic_time(:millisecond)
|
|
|
|
# Chunk size optimized for PostgreSQL work_mem=16MB
|
|
# Using smaller batches to reduce memory pressure
|
|
chunk_size = Application.get_env(:aprsme, :packet_pipeline)[:batch_size] || 100
|
|
|
|
# Use Stream for memory-efficient processing
|
|
# Process and reduce in one pass to avoid materializing the entire list
|
|
{success_count, error_count} =
|
|
packets
|
|
|> Stream.chunk_every(chunk_size)
|
|
|> Stream.map(&process_chunk/1)
|
|
|> Enum.reduce({0, 0}, fn {success, error}, {total_success, total_error} ->
|
|
{total_success + success, total_error + error}
|
|
end)
|
|
|
|
end_time = System.monotonic_time(:millisecond)
|
|
duration = end_time - start_time
|
|
|
|
# Monitor memory usage after processing (only this process)
|
|
{:memory, memory_after} = Process.info(self(), :memory)
|
|
memory_diff = memory_after - memory_before
|
|
|
|
# Get current process memory info
|
|
process_info = Process.info(self(), [:memory, :heap_size, :total_heap_size])
|
|
|
|
# Force garbage collection if memory usage is high
|
|
# 50MB threshold for memory diff (per process)
|
|
# 100MB threshold for process memory
|
|
if memory_diff > 52_428_800 or process_info[:memory] > 104_857_600 do
|
|
:erlang.garbage_collect()
|
|
|
|
Logger.warning("High memory usage detected, forced garbage collection",
|
|
memory_info:
|
|
LogSanitizer.log_data(
|
|
memory_diff_bytes: memory_diff,
|
|
process_memory: process_info[:memory],
|
|
heap_size: process_info[:heap_size],
|
|
total_heap_size: process_info[:total_heap_size],
|
|
batch_size: length(packets),
|
|
gc_forced: true
|
|
)
|
|
)
|
|
end
|
|
|
|
:telemetry.execute(
|
|
[
|
|
:aprsme,
|
|
:packet_pipeline,
|
|
:batch
|
|
],
|
|
%{
|
|
count: length(packets),
|
|
success: success_count,
|
|
error: error_count,
|
|
duration_ms: duration,
|
|
memory_diff: memory_diff
|
|
},
|
|
%{}
|
|
)
|
|
end
|
|
|
|
@spec process_chunk(list(map())) :: {non_neg_integer(), non_neg_integer()}
|
|
defp process_chunk(packets) do
|
|
# Use Stream for memory-efficient packet preparation
|
|
# Note: truncate_datetimes_to_second is already called inside prepare_packet_for_insert
|
|
# prepare_packet_for_insert returns {insert_attrs, broadcast_info} or nil
|
|
packet_stream =
|
|
packets
|
|
|> Stream.map(&prepare_packet_for_insert/1)
|
|
|> Stream.reject(&is_nil/1)
|
|
|
|
# Separate valid and invalid packets using Stream
|
|
{valid_pairs, invalid_count} =
|
|
Enum.reduce(packet_stream, {[], 0}, fn {insert_attrs, _bcast} = pair, {valid_acc, invalid_acc} ->
|
|
if valid_packet?(insert_attrs) do
|
|
{[pair | valid_acc], invalid_acc}
|
|
else
|
|
{valid_acc, invalid_acc + 1}
|
|
end
|
|
end)
|
|
|
|
# Reverse to maintain order and split into parallel lists
|
|
valid_pairs = Enum.reverse(valid_pairs)
|
|
valid_inserts = Enum.map(valid_pairs, fn {insert, _bcast} -> insert end)
|
|
|
|
# Insert valid packets in batch
|
|
# Optimized for PostgreSQL with synchronous_commit=off
|
|
insert_opts = [
|
|
# Don't return IDs for better performance
|
|
returning: false,
|
|
# Skip conflicts on the composite PK; any other error surfaces to the rescue below
|
|
on_conflict: :nothing,
|
|
conflict_target: [:id, :received_at],
|
|
# Increased timeout for large batches
|
|
timeout: 60_000
|
|
]
|
|
|
|
try do
|
|
{inserted_count, _} = Repo.insert_all(Aprsme.Packet, valid_inserts, insert_opts)
|
|
|
|
_ = broadcast_packets_async(Enum.map(valid_pairs, fn {_insert, bcast} -> bcast end))
|
|
|
|
{inserted_count, invalid_count}
|
|
rescue
|
|
error ->
|
|
Logger.error("Batch insert failed: #{inspect(error)}, falling back to individual inserts")
|
|
|
|
# Fall back to individual inserts so partial success is possible
|
|
{fallback_inserted, fallback_bcasts} = insert_individually(valid_pairs)
|
|
|
|
# Broadcast whatever was successfully inserted
|
|
_ =
|
|
if fallback_bcasts != [] do
|
|
broadcast_packets_async(fallback_bcasts)
|
|
end
|
|
|
|
{fallback_inserted, invalid_count + length(valid_pairs) - fallback_inserted}
|
|
end
|
|
end
|
|
|
|
# Fall back to inserting packets one at a time when batch insert fails.
|
|
# Returns {inserted_count, list_of_broadcast_infos_for_successes}.
|
|
@spec insert_individually(list({map(), map()})) :: {non_neg_integer(), list(map())}
|
|
defp insert_individually(pairs) do
|
|
Enum.reduce(pairs, {0, []}, fn {insert_attrs, bcast}, {count, inserted} ->
|
|
try do
|
|
changeset = Aprsme.Packet.changeset(%Aprsme.Packet{}, insert_attrs)
|
|
|
|
case Repo.insert(changeset, on_conflict: :nothing, conflict_target: [:id, :received_at]) do
|
|
{:ok, _record} ->
|
|
{count + 1, [bcast | inserted]}
|
|
|
|
{:error, _changeset} ->
|
|
{count, inserted}
|
|
end
|
|
rescue
|
|
error ->
|
|
Logger.error("Individual insert failed: #{inspect(error)}")
|
|
{count, inserted}
|
|
end
|
|
end)
|
|
end
|
|
|
|
# Broadcast packets asynchronously using supervised task pool.
|
|
# Failures here were previously silent — the Task died, the broadcast never
|
|
# reached subscribers, and there was no log trail. Rescue and log so
|
|
# operational issues (PubSub down, serialization bugs) are observable.
|
|
defp broadcast_packets_async(bcasts) do
|
|
fun = fn ->
|
|
cluster_enabled = Application.get_env(:aprsme, :cluster_enabled, false)
|
|
|
|
Enum.each(bcasts, fn bcast ->
|
|
try do
|
|
broadcast_single_packet(bcast, cluster_enabled)
|
|
rescue
|
|
error ->
|
|
Logger.error("Packet broadcast failed: #{inspect(error)}",
|
|
broadcast_error:
|
|
LogSanitizer.log_data(
|
|
error: Exception.message(error),
|
|
routing_callsign: Map.get(bcast, :routing_callsign)
|
|
)
|
|
)
|
|
end
|
|
end)
|
|
end
|
|
|
|
if test_env?() do
|
|
fun.()
|
|
{:ok, self()}
|
|
else
|
|
Aprsme.BroadcastTaskSupervisor.async_execute(fun)
|
|
end
|
|
end
|
|
|
|
defp broadcast_single_packet(
|
|
%{payload: packet, routing_callsign: routing_callsign, has_weather: has_weather?},
|
|
cluster_enabled
|
|
) do
|
|
if cluster_enabled do
|
|
PacketDistributor.distribute_packet(packet)
|
|
else
|
|
Aprsme.StreamingPacketsPubSub.broadcast_packet(packet)
|
|
Aprsme.SpatialPubSub.broadcast_packet(packet)
|
|
end
|
|
|
|
broadcast_legacy_topics(packet, routing_callsign, has_weather?)
|
|
end
|
|
|
|
# Replaces the `aprs_packets` pg_notify path (removed as a DB trigger).
|
|
defp broadcast_legacy_topics(packet, routing_callsign, has_weather?) do
|
|
# Hot path: PubSub.broadcast failures tolerated (logged by PubSub itself).
|
|
_ = Phoenix.PubSub.broadcast(Aprsme.PubSub, "postgres:aprsme_packets", {:postgres_packet, packet})
|
|
|
|
_ =
|
|
if is_binary(routing_callsign) and routing_callsign != "" do
|
|
_ = Phoenix.PubSub.broadcast(Aprsme.PubSub, "packets:#{routing_callsign}", {:postgres_packet, packet})
|
|
|
|
if has_weather? do
|
|
Phoenix.PubSub.broadcast(Aprsme.PubSub, "weather:#{routing_callsign}", {:weather_packet, packet})
|
|
end
|
|
end
|
|
|
|
:ok
|
|
end
|
|
|
|
defp test_env? do
|
|
Application.get_env(:aprsme, :env) == :test
|
|
end
|
|
|
|
defp prepare_packet_for_insert(packet_data) do
|
|
# Reuse the received_at stamped in Is.dispatch; fall back if missing.
|
|
# packet_data is already a plain map — Is.dispatch ran struct_to_map before submitting.
|
|
current_time =
|
|
case Map.get(packet_data, :received_at) do
|
|
%DateTime{} = dt -> DateTime.truncate(dt, :microsecond)
|
|
_ -> DateTime.truncate(DateTime.utc_now(), :microsecond)
|
|
end
|
|
|
|
attrs = Map.put(packet_data, :received_at, current_time)
|
|
|
|
# Extract additional data from the parsed packet including raw packet
|
|
attrs = Aprsme.Packet.extract_additional_data(attrs, attrs[:raw_packet] || "")
|
|
|
|
# Detect and set item/object fields
|
|
attrs = detect_item_or_object(attrs)
|
|
|
|
# Sanitize packet data to prevent database field overflow
|
|
attrs = Aprsme.PacketSanitizer.sanitize_packet(attrs)
|
|
|
|
# Normalize data_type to string if it's an atom
|
|
attrs = normalize_data_type(attrs)
|
|
|
|
# Apply the same processing as the original store_packet function.
|
|
# `enriched` holds everything needed for both insert and broadcast;
|
|
# we pre-build the broadcast payload once so the async broadcast path
|
|
# doesn't have to Map.drop/Map.merge/pick-identifier per packet.
|
|
enriched =
|
|
attrs
|
|
|> convert_coordinate_field_names()
|
|
|> convert_field_names()
|
|
|> normalize_packet_attrs()
|
|
|> set_received_at()
|
|
|> patch_lat_lon_from_data_extended()
|
|
|> then(fn attrs ->
|
|
{lat, lon} = extract_position(attrs)
|
|
set_lat_lon(attrs, lat, lon)
|
|
end)
|
|
|> normalize_ssid()
|
|
|> then(fn attrs ->
|
|
Map.put(attrs, :device_identifier, Aprsme.DeviceParser.extract_device_identifier(attrs))
|
|
end)
|
|
|> sanitize_packet_strings()
|
|
|> create_location_geometry()
|
|
|> Map.put(:inserted_at, current_time)
|
|
|> Map.put(:updated_at, current_time)
|
|
|> Map.put(:id, Ecto.UUID.generate())
|
|
|> Map.delete("id")
|
|
|> Map.delete(:data_extended)
|
|
|> normalize_numeric_types()
|
|
|> truncate_datetimes_to_second()
|
|
|> set_has_weather()
|
|
|
|
insert_attrs = remove_non_schema_fields(enriched)
|
|
bcast_info = build_broadcast_info(enriched)
|
|
|
|
{insert_attrs, bcast_info}
|
|
rescue
|
|
error ->
|
|
Logger.error("Failed to prepare packet for batch insert: #{inspect(error)}")
|
|
nil
|
|
end
|
|
|
|
# Pre-build the broadcast payload and routing info once per packet.
|
|
# `:sender` is replaced with the display identifier (object/item name)
|
|
# while `routing_callsign` preserves the original sender for topic routing.
|
|
defp build_broadcast_info(enriched) do
|
|
identifier =
|
|
enriched[:object_name] || enriched[:item_name] || enriched[:sender]
|
|
|
|
payload =
|
|
enriched
|
|
|> Map.drop([:location, :inserted_at, :updated_at])
|
|
|> Map.put(:sender, identifier)
|
|
|> Map.put(:latitude, enriched[:lat])
|
|
|> Map.put(:longitude, enriched[:lon])
|
|
|
|
original_sender = enriched[:sender] || enriched[:base_callsign]
|
|
|
|
routing_callsign =
|
|
if is_binary(original_sender) and original_sender != "" do
|
|
original_sender |> String.trim() |> String.upcase()
|
|
end
|
|
|
|
%{
|
|
payload: payload,
|
|
routing_callsign: routing_callsign,
|
|
has_weather: enriched[:has_weather] == true
|
|
}
|
|
end
|
|
|
|
# Compute has_weather from any populated weather field (replaces DB trigger)
|
|
defp set_has_weather(attrs) do
|
|
has_weather? =
|
|
Enum.any?(Aprsme.EncodingUtils.weather_fields(), fn field ->
|
|
not is_nil(Map.get(attrs, field))
|
|
end)
|
|
|
|
Map.put(attrs, :has_weather, has_weather?)
|
|
end
|
|
|
|
# Helper function to remove fields that exist in parser output but not in database schema
|
|
defp remove_non_schema_fields(attrs) do
|
|
# Use whitelist approach - only keep fields that are in our schema
|
|
Aprsme.PacketFieldWhitelist.filter_fields(attrs)
|
|
end
|
|
|
|
defp valid_packet?(nil), do: false
|
|
defp valid_packet?(%{sender: sender}) when is_binary(sender) and byte_size(sender) > 0, do: true
|
|
defp valid_packet?(_), do: false
|
|
|
|
# Detect if packet is an item or object and set appropriate fields
|
|
defp detect_item_or_object(attrs) do
|
|
cond do
|
|
object_packet?(attrs) -> apply_object_fields(attrs)
|
|
item_packet?(attrs) -> apply_item_fields(attrs)
|
|
true -> attrs
|
|
end
|
|
end
|
|
|
|
defp object_packet?(attrs) do
|
|
info_field = get_in(attrs, [:data, "information_field"])
|
|
|
|
attrs[:data_type] == "object" or attrs["data_type"] == "object" or
|
|
(is_binary(info_field) and String.starts_with?(info_field, ";"))
|
|
end
|
|
|
|
defp item_packet?(attrs) do
|
|
Map.has_key?(attrs, :itemname) or Map.has_key?(attrs, "itemname")
|
|
end
|
|
|
|
defp apply_object_fields(attrs) do
|
|
info_field = get_in(attrs, [:data, "information_field"])
|
|
|
|
object_name =
|
|
extract_object_name(attrs) ||
|
|
extract_object_name_from_info_field(info_field)
|
|
|
|
attrs
|
|
|> Map.put(:object_name, object_name)
|
|
|> Map.put(:is_object, true)
|
|
|> Map.delete(:itemname)
|
|
|> Map.delete("itemname")
|
|
end
|
|
|
|
defp apply_item_fields(attrs) do
|
|
item_name = Map.get(attrs, :itemname) || Map.get(attrs, "itemname")
|
|
|
|
attrs
|
|
|> Map.put(:item_name, item_name)
|
|
|> Map.put(:is_item, true)
|
|
|> Map.delete(:itemname)
|
|
|> Map.delete("itemname")
|
|
end
|
|
|
|
defp extract_object_name(attrs) do
|
|
case attrs[:data_extended] do
|
|
%{name: name} when is_binary(name) -> name
|
|
%{"name" => name} when is_binary(name) -> name
|
|
_ -> nil
|
|
end
|
|
end
|
|
|
|
defp extract_object_name_from_info_field(nil), do: nil
|
|
|
|
defp extract_object_name_from_info_field(info_field) do
|
|
# Object format: ;OBJECTNAM*DDHHMMz... or ;OBJECTNAM_DDHHMMz... (killed)
|
|
# Object names are 9 printable ASCII characters followed by * (alive) or _ (killed)
|
|
case Regex.run(~r/^;(.{9})[*_]/, info_field) do
|
|
[_, name] -> String.trim(name)
|
|
_ -> nil
|
|
end
|
|
end
|
|
|
|
# Convert field names that don't match our schema
|
|
defp convert_field_names(attrs) do
|
|
attrs
|
|
|> then(fn a ->
|
|
# Convert aprs_messaging? to aprs_messaging
|
|
case Map.get(a, :aprs_messaging?) || Map.get(a, "aprs_messaging?") do
|
|
nil ->
|
|
a
|
|
|
|
value ->
|
|
a
|
|
|> Map.put(:aprs_messaging, value)
|
|
|> Map.delete(:aprs_messaging?)
|
|
|> Map.delete("aprs_messaging?")
|
|
end
|
|
end)
|
|
|> then(fn a ->
|
|
# Convert timestamp from integer to string if needed
|
|
case Map.get(a, :timestamp) do
|
|
timestamp when is_integer(timestamp) ->
|
|
Map.put(a, :timestamp, Integer.to_string(timestamp))
|
|
|
|
_ ->
|
|
a
|
|
end
|
|
end)
|
|
end
|
|
|
|
# Convert latitude/longitude to lat/lon if present at top level
|
|
defp convert_coordinate_field_names(attrs) do
|
|
attrs
|
|
|> then(fn a ->
|
|
case Map.get(a, :latitude) do
|
|
nil -> a
|
|
lat -> a |> Map.put(:lat, lat) |> Map.delete(:latitude)
|
|
end
|
|
end)
|
|
|> then(fn a ->
|
|
case Map.get(a, :longitude) do
|
|
nil -> a
|
|
lon -> a |> Map.put(:lon, lon) |> Map.delete(:longitude)
|
|
end
|
|
end)
|
|
end
|
|
|
|
# Helper functions copied from Packets module for consistency
|
|
defp normalize_packet_attrs(attrs) do
|
|
attrs
|
|
|> Map.put_new(:base_callsign, attrs[:sender])
|
|
|> Map.put_new(:data_type, "unknown")
|
|
|> Map.put_new(:destination, "")
|
|
|> Map.put_new(:path, "")
|
|
|> Map.put_new(:ssid, "")
|
|
|> Map.put_new(:data_extended, %{})
|
|
end
|
|
|
|
defp set_received_at(attrs) do
|
|
received_at = attrs[:received_at] || DateTime.utc_now()
|
|
Map.put(attrs, :received_at, received_at)
|
|
end
|
|
|
|
defp patch_lat_lon_from_data_extended(attrs) do
|
|
case attrs[:data_extended] do
|
|
%{latitude: lat, longitude: lon} when not is_nil(lat) and not is_nil(lon) ->
|
|
attrs
|
|
|> Map.put(:lat, lat)
|
|
|> Map.put(:lon, lon)
|
|
|> Map.put(:has_position, true)
|
|
|
|
_ ->
|
|
attrs
|
|
end
|
|
end
|
|
|
|
defp extract_position(packet_data) do
|
|
if not is_nil(packet_data[:lat]) and not is_nil(packet_data[:lon]) do
|
|
{Aprsme.EncodingUtils.to_float(packet_data.lat), Aprsme.EncodingUtils.to_float(packet_data.lon)}
|
|
else
|
|
extract_position_from_data_extended(packet_data[:data_extended])
|
|
end
|
|
end
|
|
|
|
defp extract_position_from_data_extended(nil), do: {nil, nil}
|
|
|
|
defp extract_position_from_data_extended(data_extended) when is_map(data_extended) do
|
|
if has_standard_position?(data_extended) do
|
|
extract_standard_position(data_extended)
|
|
else
|
|
extract_position_from_data_extended_case(data_extended)
|
|
end
|
|
end
|
|
|
|
defp extract_position_from_data_extended(_), do: {nil, nil}
|
|
|
|
defp has_standard_position?(data_extended) when is_map(data_extended) and not is_struct(data_extended) do
|
|
not is_nil(data_extended[:latitude]) and not is_nil(data_extended[:longitude])
|
|
end
|
|
|
|
defp has_standard_position?(_), do: false
|
|
|
|
defp extract_standard_position(data_extended) when is_map(data_extended) and not is_struct(data_extended) do
|
|
{Aprsme.EncodingUtils.to_float(data_extended[:latitude]), Aprsme.EncodingUtils.to_float(data_extended[:longitude])}
|
|
end
|
|
|
|
defp extract_standard_position(_), do: {nil, nil}
|
|
|
|
defp extract_position_from_data_extended_case(data_extended) do
|
|
lat = extract_lat_from_ext_map(data_extended)
|
|
lon = extract_lon_from_ext_map(data_extended)
|
|
{Aprsme.EncodingUtils.to_float(lat), Aprsme.EncodingUtils.to_float(lon)}
|
|
end
|
|
|
|
defp extract_lat_from_ext_map(ext_map) do
|
|
ext_map[:latitude] || ext_map["latitude"] ||
|
|
(Map.has_key?(ext_map, :position) &&
|
|
(ext_map[:position][:latitude] || ext_map[:position]["latitude"])) ||
|
|
(Map.has_key?(ext_map, "position") &&
|
|
(ext_map["position"][:latitude] || ext_map["position"]["latitude"]))
|
|
end
|
|
|
|
defp extract_lon_from_ext_map(ext_map) do
|
|
ext_map[:longitude] || ext_map["longitude"] ||
|
|
(Map.has_key?(ext_map, :position) &&
|
|
(ext_map[:position][:longitude] || ext_map[:position]["longitude"])) ||
|
|
(Map.has_key?(ext_map, "position") &&
|
|
(ext_map["position"][:longitude] || ext_map["position"]["longitude"]))
|
|
end
|
|
|
|
# extract_position already ran to_float; lat/lon are floats or nil here.
|
|
# Round and store whatever we have; range validation happens in
|
|
# create_location_geometry/1 where it controls Geo.Point construction.
|
|
defp set_lat_lon(attrs, lat, lon) do
|
|
attrs
|
|
|> Map.put(:lat, round_coord(lat))
|
|
|> Map.put(:lon, round_coord(lon))
|
|
|> Map.put(:has_position, is_float(lat) and is_float(lon))
|
|
end
|
|
|
|
defp round_coord(n) when is_float(n), do: Float.round(n, 6)
|
|
defp round_coord(_), do: nil
|
|
|
|
# Build the PostGIS point only for in-range coordinates.
|
|
# Must run AFTER sanitize_packet_strings so the struct survives intact.
|
|
defp create_location_geometry(%{lat: lat, lon: lon} = attrs)
|
|
when is_float(lat) and is_float(lon) and lat >= -90.0 and lat <= 90.0 and lon >= -180.0 and lon <= 180.0 do
|
|
Map.put(attrs, :location, %Geo.Point{coordinates: {lon, lat}, srid: 4326})
|
|
end
|
|
|
|
defp create_location_geometry(attrs), do: attrs
|
|
|
|
defp normalize_ssid(attrs) do
|
|
case Map.get(attrs, :ssid) do
|
|
nil -> attrs
|
|
ssid -> Map.put(attrs, :ssid, to_string(ssid))
|
|
end
|
|
end
|
|
|
|
defp sanitize_packet_strings(value), do: Aprsme.EncodingUtils.sanitize_packet_strings(value)
|
|
|
|
defp normalize_data_type(attrs), do: Aprsme.EncodingUtils.normalize_data_type(attrs)
|
|
|
|
defp truncate_datetimes_to_second(%DateTime{} = dt), do: DateTime.truncate(dt, :second)
|
|
defp truncate_datetimes_to_second({:ok, %DateTime{} = dt}), do: DateTime.truncate(dt, :second)
|
|
defp truncate_datetimes_to_second({:error, _reason}), do: nil
|
|
|
|
defp truncate_datetimes_to_second(term) when is_map(term) and not is_struct(term) do
|
|
Map.new(term, fn {k, v} -> {k, truncate_datetimes_to_second(v)} end)
|
|
end
|
|
|
|
defp truncate_datetimes_to_second(list) when is_list(list), do: Enum.map(list, &truncate_datetimes_to_second/1)
|
|
defp truncate_datetimes_to_second(other), do: other
|
|
|
|
defp normalize_numeric_types(attrs) do
|
|
# Convert integer values to floats for float fields
|
|
float_fields = [
|
|
:temperature,
|
|
:humidity,
|
|
:wind_speed,
|
|
:wind_gust,
|
|
:pressure,
|
|
:rain_1h,
|
|
:rain_24h,
|
|
:rain_since_midnight,
|
|
:snow,
|
|
:speed,
|
|
:altitude
|
|
]
|
|
|
|
Enum.reduce(float_fields, attrs, fn field, acc ->
|
|
case Map.get(acc, field) do
|
|
value when is_integer(value) -> Map.put(acc, field, value * 1.0)
|
|
_ -> acc
|
|
end
|
|
end)
|
|
end
|
|
end
|