From 0beb278371444395b1bafccf60e7c8a730e29374 Mon Sep 17 00:00:00 2001
From: Graham McIntire
Date: Tue, 24 Jun 2025 15:06:09 -0500
Subject: [PATCH] add packet batcher
---
config/config.exs | 9 +-
lib/aprsme/application.ex | 6 +-
lib/aprsme/data_extended.ex | 16 +-
lib/aprsme/is/is.ex | 67 ++--
lib/aprsme/packet.ex | 3 +-
lib/aprsme/packet_consumer.ex | 369 +++++++++++++++++++
lib/aprsme/packet_pipeline_setup.ex | 40 ++
lib/aprsme/packet_pipeline_supervisor.ex | 22 ++
lib/aprsme/packet_producer.ex | 56 +++
lib/aprsme/packets.ex | 25 +-
lib/aprsme_web/components/core_components.ex | 34 +-
lib/aprsme_web/live/api_docs_live.ex | 18 +-
lib/aprsme_web/live/map_live/index.ex | 16 +
lib/aprsme_web/telemetry.ex | 11 +-
mix.exs | 1 +
mix.lock | 1 +
test/aprsme/packets_mic_e_test.exs | 34 ++
test/support/conn_case.ex | 1 +
18 files changed, 659 insertions(+), 70 deletions(-)
create mode 100644 lib/aprsme/packet_consumer.ex
create mode 100644 lib/aprsme/packet_pipeline_setup.ex
create mode 100644 lib/aprsme/packet_pipeline_supervisor.ex
create mode 100644 lib/aprsme/packet_producer.ex
diff --git a/config/config.exs b/config/config.exs
index 56b89b8..01dd85c 100644
--- a/config/config.exs
+++ b/config/config.exs
@@ -50,7 +50,14 @@ config :aprsme,
aprs_is_login_id: System.get_env("APRS_CALLSIGN"),
aprs_is_password: System.get_env("APRS_PASSCODE"),
auto_migrate: true,
- env: config_env()
+ env: config_env(),
+ # GenStage packet processing configuration
+ packet_pipeline: [
+ max_buffer_size: 1000,
+ batch_size: 100,
+ batch_timeout: 1000,
+ max_demand: 50
+ ]
# Configure esbuild (the version is required)
config :esbuild,
diff --git a/lib/aprsme/application.ex b/lib/aprsme/application.ex
index 8baa709..e6305d3 100644
--- a/lib/aprsme/application.ex
+++ b/lib/aprsme/application.ex
@@ -33,7 +33,11 @@ defmodule Aprsme.Application do
{Oban, :aprsme |> Application.get_env(Oban, []) |> Keyword.put(:queues, default: 10, maintenance: 2)},
Aprsme.Presence,
Aprsme.AprsIsConnection,
- Aprsme.PostgresNotifier
+ Aprsme.PostgresNotifier,
+ # Start the packet processing pipeline
+ Aprsme.PacketPipelineSupervisor,
+ # Start the packet pipeline setup
+ Aprsme.PacketPipelineSetup
]
children = maybe_add_is_supervisor(children, Application.get_env(:aprsme, :env))
diff --git a/lib/aprsme/data_extended.ex b/lib/aprsme/data_extended.ex
index 0eac347..3650e5d 100644
--- a/lib/aprsme/data_extended.ex
+++ b/lib/aprsme/data_extended.ex
@@ -14,6 +14,13 @@ defmodule Aprsme.DataExtended do
field :longitude, :decimal
field :symbol_code, :string
field :symbol_table_id, :string
+ field :timestamp, :string
+ field :speed, :float
+ field :course, :integer
+ field :has_position, :boolean, default: false
+ field :compressed, :boolean, default: false
+ field :dao, :string
+ field :position_ambiguity, :integer, default: 0
end
@type t :: %__MODULE__{}
@@ -29,7 +36,14 @@ defmodule Aprsme.DataExtended do
:latitude,
:longitude,
:symbol_code,
- :symbol_table_id
+ :symbol_table_id,
+ :timestamp,
+ :speed,
+ :course,
+ :has_position,
+ :compressed,
+ :dao,
+ :position_ambiguity
])
|> validate_required([
:aprs_messaging,
diff --git a/lib/aprsme/is/is.ex b/lib/aprsme/is/is.ex
index 6b9790e..958d657 100644
--- a/lib/aprsme/is/is.ex
+++ b/lib/aprsme/is/is.ex
@@ -53,11 +53,11 @@ defmodule Aprsme.Is do
Process.sleep(2000)
# Get startup parameters
- server = Application.get_env(:aprsme, :aprsme_is_server, ~c"rotate.aprs2.net")
- port = Application.get_env(:aprsme, :aprsme_is_port, 14_580)
- default_filter = Application.get_env(:aprsme, :aprsme_is_default_filter, "r/33/-96/100")
- aprs_user_id = Application.get_env(:aprsme, :aprsme_is_login_id, "W5ISP")
- aprs_passcode = Application.get_env(:aprsme, :aprsme_is_password, "-1")
+ server = Application.get_env(:aprsme, :aprs_is_server, ~c"dallas.aprs2.net")
+ 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()
@@ -412,48 +412,31 @@ defmodule Aprsme.Is do
case Aprs.parse(message) do
{:ok, parsed_message} ->
# Store the packet in the database for future replay
- # Use Task to avoid slowing down the main process
- Task.start(fn ->
- require Logger
+ # Use GenStage pipeline for efficient batch processing
+ try do
+ # Always set received_at timestamp to ensure consistency
+ current_time = DateTime.truncate(DateTime.utc_now(), :microsecond)
+ packet_data = Map.put(parsed_message, :received_at, current_time)
- try do
- # Always set received_at timestamp to ensure consistency
- current_time = DateTime.truncate(DateTime.utc_now(), :microsecond)
- packet_data = Map.put(parsed_message, :received_at, current_time)
+ # Convert to map before storing to avoid struct conversion issues
+ attrs = struct_to_map(packet_data)
- # Convert to map before storing to avoid struct conversion issues
- attrs = struct_to_map(packet_data)
+ # Extract additional data from the parsed packet including raw packet
+ attrs = Aprsme.Packet.extract_additional_data(attrs, message)
- # Extract additional data from the parsed packet including raw packet
- attrs = Aprsme.Packet.extract_additional_data(attrs, message)
+ # Normalize data_type to string if it's an atom
+ attrs = normalize_data_type(attrs)
- # Normalize data_type to string if it's an atom
- attrs = normalize_data_type(attrs)
+ # Submit to GenStage pipeline for batch processing
+ Aprsme.PacketProducer.submit_packet(attrs)
+ rescue
+ error ->
+ require Logger
- # Store in database through the Packets context
- case Aprsme.Packets.store_packet(attrs) do
- {:ok, _packet} ->
- # Packet stored successfully
- :ok
-
- {:error, :storage_exception} ->
- Logger.error("Storage exception while storing packet from #{inspect(parsed_message.sender)}")
-
- Logger.debug("Packet attributes that failed: #{inspect(attrs)}")
-
- {:error, :validation_error} ->
- Logger.error("Validation error while storing packet from #{inspect(parsed_message.sender)}")
-
- Logger.debug("Packet attributes that failed: #{inspect(attrs)}")
- end
- rescue
- error ->
- Logger.error("Exception while storing packet from #{inspect(parsed_message.sender)}: #{inspect(error)}")
-
- Logger.debug("Raw message: #{inspect(message)}")
- Logger.debug("Parsed message: #{inspect(parsed_message)}")
- end
- end)
+ 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
# Broadcast to live clients
AprsmeWeb.Endpoint.broadcast("aprs_messages", "packet", parsed_message)
diff --git a/lib/aprsme/packet.ex b/lib/aprsme/packet.ex
index 9464a26..66f5e8f 100644
--- a/lib/aprsme/packet.ex
+++ b/lib/aprsme/packet.ex
@@ -147,7 +147,8 @@ defmodule Aprsme.Packet do
defp extract_coordinates_from_data_extended(nil), do: {nil, nil}
- defp extract_coordinates_from_data_extended(data_extended) when is_map(data_extended) do
+ defp extract_coordinates_from_data_extended(data_extended)
+ when is_map(data_extended) and not is_struct(data_extended) do
{data_extended[:latitude], data_extended[:longitude]}
end
diff --git a/lib/aprsme/packet_consumer.ex b/lib/aprsme/packet_consumer.ex
new file mode 100644
index 0000000..343b37d
--- /dev/null
+++ b/lib/aprsme/packet_consumer.ex
@@ -0,0 +1,369 @@
+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.Repo
+
+ require Logger
+
+ def start_link(opts \\ []) do
+ GenStage.start_link(__MODULE__, opts, name: __MODULE__)
+ end
+
+ @impl true
+ def init(opts) do
+ batch_size = opts[:batch_size] || 100
+ batch_timeout = opts[:batch_timeout] || 1000
+
+ # Start a timer for batch processing
+ timer = Process.send_after(self(), :process_batch, batch_timeout)
+
+ {:consumer,
+ %{
+ batch: [],
+ batch_size: batch_size,
+ batch_timeout: batch_timeout,
+ timer: timer
+ }}
+ end
+
+ @impl true
+ def handle_events(events, _from, %{batch: batch, batch_size: batch_size} = state) do
+ new_batch = batch ++ events
+
+ if length(new_batch) >= batch_size do
+ # Process the batch immediately
+ process_batch(new_batch)
+ {:noreply, [], %{state | batch: []}}
+ else
+ # Add to batch and wait for more
+ {:noreply, [], %{state | batch: new_batch}}
+ end
+ end
+
+ @impl true
+ def handle_info(:process_batch, %{batch: batch, batch_timeout: timeout} = state) do
+ if length(batch) > 0 do
+ process_batch(batch)
+ end
+
+ # Start a new timer
+ timer = Process.send_after(self(), :process_batch, timeout)
+ {:noreply, [], %{state | batch: [], timer: timer}}
+ end
+
+ defp process_batch(packets) do
+ require Logger
+
+ start_time = System.monotonic_time(:millisecond)
+
+ results =
+ packets
+ |> Enum.chunk_every(50)
+ |> Enum.map(&process_chunk/1)
+
+ {success_count, error_count} =
+ Enum.reduce(results, {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
+
+ :telemetry.execute(
+ [
+ :aprsme,
+ :packet_pipeline,
+ :batch
+ ],
+ %{count: length(packets), success: success_count, error: error_count, duration_ms: duration},
+ %{}
+ )
+
+ Logger.info(
+ "Processed batch of #{length(packets)} packets in #{duration}ms (success: #{success_count}, errors: #{error_count})"
+ )
+ end
+
+ defp process_chunk(packets) do
+ # Prepare packets for batch insertion
+ packet_attrs =
+ packets
+ |> Enum.map(&prepare_packet_for_insert/1)
+ # Ensure truncation here
+ |> Enum.map(&truncate_datetimes_to_second/1)
+
+ # Filter out invalid packets
+ {valid_packets, invalid_packets} = Enum.split_with(packet_attrs, &valid_packet?/1)
+
+ # Insert valid packets in batch
+ case Repo.insert_all(Aprsme.Packet, valid_packets, returning: [:id]) do
+ {:error, error} ->
+ Logger.error("Batch insert failed: #{inspect(error)}")
+ {0, length(packets)}
+
+ {inserted_count, _} ->
+ error_count = Enum.count(invalid_packets)
+ {inserted_count, error_count}
+ end
+ end
+
+ defp prepare_packet_for_insert(packet_data) do
+ # Always set received_at timestamp to ensure consistency
+ current_time = DateTime.truncate(DateTime.utc_now(), :microsecond)
+ packet_data = Map.put(packet_data, :received_at, current_time)
+
+ # Convert to map before storing to avoid struct conversion issues
+ attrs = struct_to_map(packet_data)
+
+ # Extract additional data from the parsed packet including raw packet
+ attrs = Aprsme.Packet.extract_additional_data(attrs, attrs[:raw_packet] || "")
+
+ # 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
+ attrs
+ |> 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 ->
+ device_identifier = Aprsme.DeviceParser.extract_device_identifier(packet_data)
+ matched_device = Aprsme.DeviceIdentification.lookup_device_by_identifier(device_identifier)
+ canonical_identifier = if matched_device, do: matched_device.identifier, else: device_identifier
+ Map.put(attrs, :device_identifier, canonical_identifier)
+ end)
+ |> sanitize_packet_strings()
+ |> Map.put(:inserted_at, current_time)
+ |> Map.put(:updated_at, current_time)
+ |> Map.delete(:id)
+ |> Map.delete("id")
+ # Remove embedded field for batch insert
+ |> Map.delete(:data_extended)
+ |> normalize_numeric_types()
+ |> sanitize_raw_packet()
+ |> then(fn attrs -> Map.new(attrs, fn {k, v} -> {k, sanitize_packet_strings(v)} end) end)
+ |> truncate_datetimes_to_second()
+ rescue
+ error ->
+ Logger.error("Failed to prepare packet for batch insert: #{inspect(error)}")
+ nil
+ 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
+
+ # 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(:information_field, "")
+ |> 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
+ {to_float(packet_data.lat), 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
+ {to_float(data_extended[:latitude]), 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)
+ {to_float(lat), 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
+
+ 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)
+
+ n when is_integer(n) ->
+ n * 1.0
+
+ n when is_binary(n) ->
+ case Float.parse(n) do
+ {f, _} -> Float.round(f, 6)
+ :error -> nil
+ end
+
+ _ ->
+ nil
+ end
+
+ attrs
+ |> Map.put(:lat, round6.(lat))
+ |> Map.put(:lon, round6.(lon))
+ |> Map.put(:has_position, not is_nil(lat) and not is_nil(lon))
+ end
+
+ 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(%DateTime{} = dt), do: dt
+ defp sanitize_packet_strings(%NaiveDateTime{} = ndt), do: ndt
+ defp sanitize_packet_strings(%_struct{} = struct), do: struct |> Map.from_struct() |> sanitize_packet_strings()
+ defp sanitize_packet_strings(list) when is_list(list), do: Enum.map(list, &sanitize_packet_strings/1)
+
+ defp sanitize_packet_strings(binary) when is_binary(binary) do
+ s = Aprsme.EncodingUtils.sanitize_string(binary)
+ if is_binary(s), do: s, else: ""
+ end
+
+ defp sanitize_packet_strings(other), do: other
+
+ defp to_float(value) when is_float(value), do: value
+ defp to_float(value) when is_integer(value), do: value * 1.0
+ defp to_float(%Decimal{} = value), do: Decimal.to_float(value)
+
+ defp to_float(value) when is_binary(value) do
+ case Float.parse(value) do
+ {float, _} -> float
+ :error -> nil
+ end
+ end
+
+ defp to_float(_), do: nil
+
+ defp normalize_data_type(%{data_type: data_type} = attrs) when is_atom(data_type) do
+ %{attrs | data_type: to_string(data_type)}
+ end
+
+ defp normalize_data_type(attrs), do: attrs
+
+ 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)
+
+ 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 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(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,
+ :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
+
+ defp sanitize_raw_packet(attrs) do
+ # Use the exact same sanitization approach as the original working code
+ sanitize_packet_strings(attrs)
+ end
+end
diff --git a/lib/aprsme/packet_pipeline_setup.ex b/lib/aprsme/packet_pipeline_setup.ex
new file mode 100644
index 0000000..10d6d5c
--- /dev/null
+++ b/lib/aprsme/packet_pipeline_setup.ex
@@ -0,0 +1,40 @@
+defmodule Aprsme.PacketPipelineSetup do
+ @moduledoc """
+ Handles the setup of the GenStage pipeline subscription after the supervisor starts.
+ """
+ use GenServer
+
+ def start_link(opts \\ []) do
+ GenServer.start_link(__MODULE__, opts, name: __MODULE__)
+ end
+
+ @impl true
+ def init(_opts) do
+ # Wait a moment for the pipeline to start, then set up the subscription
+ Process.send_after(self(), :setup_subscription, 100)
+ {:ok, %{}}
+ end
+
+ @impl true
+ def handle_info(:setup_subscription, state) do
+ # Set up the subscription between producer and consumer
+ config = Application.get_env(:aprsme, :packet_pipeline, [])
+ max_demand = config[:max_demand] || 50
+
+ case GenStage.sync_subscribe(Aprsme.PacketConsumer, to: Aprsme.PacketProducer, max_demand: max_demand) do
+ {:ok, _subscription} ->
+ require Logger
+
+ Logger.info("GenStage packet pipeline subscription established")
+ {:noreply, state}
+
+ {:error, reason} ->
+ require Logger
+
+ Logger.error("Failed to establish GenStage subscription: #{inspect(reason)}")
+ # Retry after a delay
+ Process.send_after(self(), :setup_subscription, 1000)
+ {:noreply, state}
+ end
+ end
+end
diff --git a/lib/aprsme/packet_pipeline_supervisor.ex b/lib/aprsme/packet_pipeline_supervisor.ex
new file mode 100644
index 0000000..afb6316
--- /dev/null
+++ b/lib/aprsme/packet_pipeline_supervisor.ex
@@ -0,0 +1,22 @@
+defmodule Aprsme.PacketPipelineSupervisor do
+ @moduledoc """
+ Supervisor for the GenStage pipeline that handles packet processing.
+ """
+ use Supervisor
+
+ def start_link(opts \\ []) do
+ Supervisor.start_link(__MODULE__, opts, name: __MODULE__)
+ end
+
+ @impl true
+ def init(_opts) do
+ config = Application.get_env(:aprsme, :packet_pipeline, [])
+
+ children = [
+ {Aprsme.PacketProducer, max_buffer_size: config[:max_buffer_size] || 1000},
+ {Aprsme.PacketConsumer, batch_size: config[:batch_size] || 100, batch_timeout: config[:batch_timeout] || 1000}
+ ]
+
+ Supervisor.init(children, strategy: :one_for_one)
+ end
+end
diff --git a/lib/aprsme/packet_producer.ex b/lib/aprsme/packet_producer.ex
new file mode 100644
index 0000000..f8ce11c
--- /dev/null
+++ b/lib/aprsme/packet_producer.ex
@@ -0,0 +1,56 @@
+defmodule Aprsme.PacketProducer do
+ @moduledoc """
+ GenStage producer that handles incoming APRS packets and sends them to consumers
+ for efficient batch processing.
+ """
+ use GenStage
+
+ require Logger
+
+ def start_link(opts \\ []) do
+ GenStage.start_link(__MODULE__, opts, name: __MODULE__)
+ end
+
+ def submit_packet(packet_data) do
+ GenStage.cast(__MODULE__, {:packet, packet_data})
+ end
+
+ @impl true
+ def init(opts) do
+ {:producer, %{demand: 0, buffer: [], max_buffer_size: opts[:max_buffer_size] || 1000}}
+ end
+
+ @impl true
+ def handle_demand(incoming_demand, %{demand: demand, buffer: buffer} = state) do
+ {events, remaining_buffer, remaining_demand} = dispatch_events(buffer, demand + incoming_demand)
+ {:noreply, events, %{state | demand: remaining_demand, buffer: remaining_buffer}}
+ end
+
+ @impl true
+ def handle_cast({:packet, packet_data}, %{demand: demand, buffer: buffer, max_buffer_size: max_size} = state) do
+ if demand > 0 do
+ # We have demand, send immediately
+ {:noreply, [packet_data], %{state | demand: demand - 1}}
+ else
+ # No demand, buffer the packet
+ new_buffer = [packet_data | buffer]
+
+ if length(new_buffer) > max_size do
+ # Buffer is full, drop oldest packet
+ Logger.warning("Packet buffer full, dropping oldest packet")
+ {:noreply, [], %{state | buffer: Enum.take(new_buffer, max_size)}}
+ else
+ {:noreply, [], %{state | buffer: new_buffer}}
+ end
+ end
+ end
+
+ defp dispatch_events(buffer, demand) when demand > 0 and length(buffer) > 0 do
+ {events, remaining} = Enum.split(buffer, demand)
+ {events, remaining, demand - length(events)}
+ end
+
+ defp dispatch_events(buffer, demand) do
+ {[], buffer, demand}
+ end
+end
diff --git a/lib/aprsme/packets.ex b/lib/aprsme/packets.ex
index a97d6bb..dac536c 100644
--- a/lib/aprsme/packets.ex
+++ b/lib/aprsme/packets.ex
@@ -108,7 +108,7 @@ defmodule Aprsme.Packets do
end
end
- defp extract_lat_from_ext_map(ext_map) do
+ defp extract_lat_from_ext_map(ext_map) when is_map(ext_map) and not is_struct(ext_map) do
ext_map[:latitude] || ext_map["latitude"] ||
(Map.has_key?(ext_map, :position) &&
(ext_map[:position][:latitude] || ext_map[:position]["latitude"])) ||
@@ -116,7 +116,9 @@ defmodule Aprsme.Packets do
(ext_map["position"][:latitude] || ext_map["position"]["latitude"]))
end
- defp extract_lon_from_ext_map(ext_map) do
+ defp extract_lat_from_ext_map(_), do: nil
+
+ defp extract_lon_from_ext_map(ext_map) when is_map(ext_map) and not is_struct(ext_map) do
ext_map[:longitude] || ext_map["longitude"] ||
(Map.has_key?(ext_map, :position) &&
(ext_map[:position][:longitude] || ext_map[:position]["longitude"])) ||
@@ -124,6 +126,8 @@ defmodule Aprsme.Packets do
(ext_map["position"][:longitude] || ext_map["position"]["longitude"]))
end
+ defp extract_lon_from_ext_map(_), do: nil
+
defp set_lat_lon(attrs, lat, lon) do
round6 = fn
nil ->
@@ -230,25 +234,34 @@ defmodule Aprsme.Packets do
defp extract_position_from_data_extended(_), do: {nil, nil}
- defp has_standard_position?(data_extended) do
+ 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 extract_standard_position(data_extended) do
+ defp has_standard_position?(_), do: false
+
+ defp extract_standard_position(data_extended) when is_map(data_extended) and not is_struct(data_extended) do
{to_float(data_extended[:latitude]), to_float(data_extended[:longitude])}
end
+ defp extract_standard_position(_), do: {nil, nil}
+
defp extract_position_from_data_extended_case(data_extended) do
case data_extended do
%{__struct__: MicE} = mic_e ->
extract_position_from_mic_e_struct(mic_e)
+ %{__struct__: Aprs.Types.ParseError} ->
+ # ParseError structs don't contain position data and don't implement Access behavior
+ {nil, nil}
+
_ ->
extract_position_from_mic_e(data_extended)
end
end
- defp extract_position_from_mic_e_struct(mic_e) do
+ defp extract_position_from_mic_e_struct(%{__struct__: MicE} = mic_e) do
+ # Use Access behavior for MicE struct which implements it
lat = mic_e[:latitude]
lon = mic_e[:longitude]
@@ -259,6 +272,8 @@ defmodule Aprsme.Packets do
end
end
+ defp extract_position_from_mic_e_struct(_), do: {nil, nil}
+
defp extract_position_from_mic_e(data_extended) do
if is_number(data_extended[:lat_degrees]) and is_number(data_extended[:lat_minutes]) and
is_number(data_extended[:lon_degrees]) and is_number(data_extended[:lon_minutes]) do
diff --git a/lib/aprsme_web/components/core_components.ex b/lib/aprsme_web/components/core_components.ex
index f4e6c77..3b55b87 100644
--- a/lib/aprsme_web/components/core_components.ex
+++ b/lib/aprsme_web/components/core_components.ex
@@ -409,15 +409,7 @@ defmodule AprsmeWeb.CoreComponents do
aprs.me
-
+ <.navigation variant={:horizontal} />
"""
@@ -636,4 +628,28 @@ defmodule AprsmeWeb.CoreComponents do
defp input_equals?(val1, val2) do
Phoenix.HTML.html_escape(val1) == Phoenix.HTML.html_escape(val2)
end
+
+ @doc """
+ Renders a navigation menu component that can be used in both regular pages and map sidebars.
+ """
+ attr :class, :string, default: ""
+ attr :variant, :atom, values: [:horizontal, :vertical], default: :horizontal
+
+ def navigation(assigns) do
+ ~H"""
+
+ Home
+
+ API
+
+
+ About
+
+
+ """
+ end
end
diff --git a/lib/aprsme_web/live/api_docs_live.ex b/lib/aprsme_web/live/api_docs_live.ex
index 5a19e00..01c9ec6 100644
--- a/lib/aprsme_web/live/api_docs_live.ex
+++ b/lib/aprsme_web/live/api_docs_live.ex
@@ -224,7 +224,7 @@ defmodule AprsmeWeb.ApiDocsLive do
RESTful JSON API for accessing APRS packet data and station information.
-
+
@@ -263,7 +263,7 @@ defmodule AprsmeWeb.ApiDocsLive do
-
+
@@ -407,7 +407,7 @@ defmodule AprsmeWeb.ApiDocsLive do
-
+
@@ -527,7 +527,7 @@ defmodule AprsmeWeb.ApiDocsLive do
-
+
@@ -591,7 +591,7 @@ defmodule AprsmeWeb.ApiDocsLive do
-
+
@@ -638,7 +638,7 @@ defmodule AprsmeWeb.ApiDocsLive do
-
+
@@ -703,7 +703,7 @@ defmodule AprsmeWeb.ApiDocsLive do
-
+
<%= if @error do %>
@@ -723,7 +723,7 @@ defmodule AprsmeWeb.ApiDocsLive do
<% end %>
-
+
<%= if @api_result do %>
@@ -750,7 +750,7 @@ defmodule AprsmeWeb.ApiDocsLive do
-
+
diff --git a/lib/aprsme_web/live/map_live/index.ex b/lib/aprsme_web/live/map_live/index.ex
index 9bc10fb..268ad7a 100644
--- a/lib/aprsme_web/live/map_live/index.ex
+++ b/lib/aprsme_web/live/map_live/index.ex
@@ -906,6 +906,22 @@ defmodule AprsmeWeb.MapLive.Index do
+
+
+
+ <.navigation variant={:vertical} class="text-sm" />
+
+
diff --git a/lib/aprsme_web/telemetry.ex b/lib/aprsme_web/telemetry.ex
index 2f21028..408ed9d 100644
--- a/lib/aprsme_web/telemetry.ex
+++ b/lib/aprsme_web/telemetry.ex
@@ -79,7 +79,16 @@ defmodule AprsmeWeb.Telemetry do
summary("vm.memory.total", unit: {:byte, :kilobyte}),
summary("vm.total_run_queue_lengths.total"),
summary("vm.total_run_queue_lengths.cpu"),
- summary("vm.total_run_queue_lengths.io")
+ summary("vm.total_run_queue_lengths.io"),
+
+ # GenStage Packet Pipeline Metrics
+ summary("aprsme.packet_pipeline.batch.count", unit: :event, description: "Number of packets in each batch"),
+ summary("aprsme.packet_pipeline.batch.success",
+ unit: :event,
+ description: "Number of successful inserts per batch"
+ ),
+ summary("aprsme.packet_pipeline.batch.error", unit: :event, description: "Number of errors per batch"),
+ summary("aprsme.packet_pipeline.batch.duration_ms", unit: :millisecond, description: "Batch insert duration (ms)")
]
end
diff --git a/mix.exs b/mix.exs
index 921dcaa..5b35133 100644
--- a/mix.exs
+++ b/mix.exs
@@ -54,6 +54,7 @@ defmodule Aprsme.MixProject do
{:geo, "~> 4.0.1"},
{:geo_postgis, "~> 3.4"},
{:geocalc, "~> 0.8"},
+ {:gen_stage, "~> 1.2"},
{:gettext, "~> 0.26.2"},
{:hackney, "~> 1.24"},
{:heroicons, "~> 0.5"},
diff --git a/mix.lock b/mix.lock
index 15eaa90..77b41ac 100644
--- a/mix.lock
+++ b/mix.lock
@@ -31,6 +31,7 @@
"file_system": {:hex, :file_system, "1.1.0", "08d232062284546c6c34426997dd7ef6ec9f8bbd090eb91780283c9016840e8f", [:mix], [], "hexpm", "bfcf81244f416871f2a2e15c1b515287faa5db9c6bcf290222206d120b3d43f6"},
"finch": {:hex, :finch, "0.19.0", "c644641491ea854fc5c1bbaef36bfc764e3f08e7185e1f084e35e0672241b76d", [:mix], [{:mime, "~> 1.0 or ~> 2.0", [hex: :mime, repo: "hexpm", optional: false]}, {:mint, "~> 1.6.2 or ~> 1.7", [hex: :mint, repo: "hexpm", optional: false]}, {:nimble_options, "~> 0.4 or ~> 1.0", [hex: :nimble_options, repo: "hexpm", optional: false]}, {:nimble_pool, "~> 1.1", [hex: :nimble_pool, repo: "hexpm", optional: false]}, {:telemetry, "~> 0.4 or ~> 1.0", [hex: :telemetry, repo: "hexpm", optional: false]}], "hexpm", "fc5324ce209125d1e2fa0fcd2634601c52a787aff1cd33ee833664a5af4ea2b6"},
"floki": {:hex, :floki, "0.38.0", "62b642386fa3f2f90713f6e231da0fa3256e41ef1089f83b6ceac7a3fd3abf33", [:mix], [], "hexpm", "a5943ee91e93fb2d635b612caf5508e36d37548e84928463ef9dd986f0d1abd9"},
+ "gen_stage": {:hex, :gen_stage, "1.3.0", "d3e102a781e30de7bc14f5bacf31c4563f5390733d6d811dd020421803cae3eb", [:mix], [], "hexpm", "02b4263f8937a5039db34637879c3332dd5eb909b463493a0e7ab996dfcf19fa"},
"geo": {:hex, :geo, "4.0.1", "f4ae3fd912b0536bfe9ec3bce15eb554197ab0739c01297c8534c20dcedd561c", [:mix], [{:jason, "~> 1.4", [hex: :jason, repo: "hexpm", optional: true]}], "hexpm", "32eb624feff75d043bbdd43f67e3869c5fc729e221333271b07cdc98ba98563d"},
"geo_postgis": {:hex, :geo_postgis, "3.7.1", "614f25b42334a615bd54bb09c22030b1aac7bac8f829bd823ab1faccf093a324", [:mix], [{:ecto, "~> 3.0", [hex: :ecto, repo: "hexpm", optional: true]}, {:geo, "~> 3.6 or ~> 4.0", [hex: :geo, repo: "hexpm", optional: false]}, {:jason, "~> 1.2", [hex: :jason, repo: "hexpm", optional: true]}, {:poison, "~> 2.2 or ~> 3.0 or ~> 4.0 or ~> 5.0 or ~> 6.0", [hex: :poison, repo: "hexpm", optional: true]}, {:postgrex, ">= 0.0.0", [hex: :postgrex, repo: "hexpm", optional: false]}], "hexpm", "c20d823c600d35b7fe9ddd5be03052bb7136c57d6f1775dbd46871545e405280"},
"geocalc": {:hex, :geocalc, "0.8.5", "b9886679e44c323e5b72dcd90a64f834d775d2600af0b656ea9f07ccdacaa5a6", [:mix], [{:decimal, "~> 2.0", [hex: :decimal, repo: "hexpm", optional: false]}], "hexpm", "3870c25c78513ec0456b69324c2be1af2202961002e81fb659559e3db162c802"},
diff --git a/test/aprsme/packets_mic_e_test.exs b/test/aprsme/packets_mic_e_test.exs
index 869f513..ed56c01 100644
--- a/test/aprsme/packets_mic_e_test.exs
+++ b/test/aprsme/packets_mic_e_test.exs
@@ -210,5 +210,39 @@ defmodule Aprsme.PacketsMicETest do
# The fix prevented the KeyError - that's what we're testing
# Invalid coordinates are handled gracefully by the validation logic
end
+
+ test "handles ParseError struct in data_extended gracefully" do
+ # This test covers the specific error case from the user's report
+ # where data_extended contains a ParseError struct instead of MicE
+ parse_error_data = %Aprs.Types.ParseError{
+ error_code: :not_implemented,
+ error_message: "PHG/DFS parsing not yet implemented",
+ raw_data: nil
+ }
+
+ packet_data = %{
+ base_callsign: "LU9DCE",
+ ssid: "4",
+ sender: "LU9DCE-4",
+ destination: "ID",
+ data_type: "phg_data",
+ path: "qAR,N3HYM-9",
+ information_field: "###################",
+ data_extended: parse_error_data
+ }
+
+ # This should NOT raise an UndefinedFunctionError for ParseError.fetch/2
+ # The fix ensures we check the struct type before using Access behavior
+ assert {:ok, stored_packet} = Packets.store_packet(packet_data)
+
+ # Verify the packet was stored correctly
+ assert stored_packet.sender == "LU9DCE-4"
+ assert stored_packet.data_type == "phg_data"
+
+ # Should not have position data since ParseError doesn't contain coordinates
+ assert stored_packet.has_position != true || stored_packet.has_position == nil
+ assert is_nil(stored_packet.lat) || stored_packet.lat == nil
+ assert is_nil(stored_packet.lon) || stored_packet.lon == nil
+ end
end
end
diff --git a/test/support/conn_case.ex b/test/support/conn_case.ex
index 0b2919f..926d065 100644
--- a/test/support/conn_case.ex
+++ b/test/support/conn_case.ex
@@ -26,6 +26,7 @@ defmodule AprsmeWeb.ConnCase do
import AprsmeWeb.ConnCase
import Phoenix.ConnTest
import Plug.Conn
+
alias Aprsme.Repo
# The default endpoint for testing