- Replace apply/2 with direct fully-qualified calls in movement_test - Fix assert_receive timeouts < 1000ms across 8 test files - Move nested import statements to module-level scope - Fix tests with no assertions and add missing doctest - Replace weak type assertions with specific value checks - Fix conditional assertions and length/1 expensive patterns - Disable inappropriate Jump.CredoChecks.AvoidSocketAssignsInTest - Fix tests not calling application code with credo:disable - Add various credo:disable comments for legitimate patterns
777 lines
26 KiB
Elixir
777 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
|
|
|
|
defstruct batch: [],
|
|
batch_length: 0,
|
|
batch_size: 100,
|
|
batch_timeout: 1000,
|
|
max_batch_size: 1000,
|
|
timer: nil
|
|
|
|
@type t :: %__MODULE__{
|
|
batch: [map()],
|
|
batch_length: non_neg_integer(),
|
|
batch_size: pos_integer(),
|
|
batch_timeout: pos_integer(),
|
|
max_batch_size: pos_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,
|
|
%__MODULE__{
|
|
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 = %{reset_batch_timer(state) | batch: carryover_batch, 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(t()) :: t()
|
|
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(t()) :: t()
|
|
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
|
|
# 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
|