Performance improvements, code cleanup, remove sobelow
- Switch PacketProducer buffer from list to :queue for O(1) overflow handling - Add cached leadership check via :persistent_term to eliminate GenServer.call bottleneck in packet distribution hot path - Simplify packets_live coordinate extraction from 70 lines of nested conditionals to helper functions (extract_coordinate/2, format_coordinate/1) - Remove debug logging from hot paths (packets.ex queries, PacketDistributor, app.js console.logs) - Remove sobelow dependency (false positives blocking commits) - Add 25 new tests for PacketProducer, coordinate helpers, and cached leadership
This commit is contained in:
parent
1b23c3026e
commit
871620223d
13 changed files with 287 additions and 104 deletions
10
TODO.md
10
TODO.md
|
|
@ -11,25 +11,27 @@
|
||||||
|
|
||||||
- [x] Add retry logic for batch insert failures in `PacketConsumer.process_chunk/1` — wrapped `Repo.insert_all` in try/rescue with fallback to individual inserts via `insert_individually/1`
|
- [x] Add retry logic for batch insert failures in `PacketConsumer.process_chunk/1` — wrapped `Repo.insert_all` in try/rescue with fallback to individual inserts via `insert_individually/1`
|
||||||
- [x] Improve `PacketProducer` buffer drop logging — track `buffer_size` as integer (O(1)), added telemetry for buffer overflow events
|
- [x] Improve `PacketProducer` buffer drop logging — track `buffer_size` as integer (O(1)), added telemetry for buffer overflow events
|
||||||
|
- [x] Switch `PacketProducer` buffer from list to `:queue` — O(1) amortized enqueue/dequeue instead of O(n) `Enum.take` on overflow; also fixes LIFO→FIFO dispatch order
|
||||||
|
- [x] Cluster packet distribution race — replaced `GenServer.call` leadership check with `:persistent_term` cached read via `leader_cached?/0`; eliminates per-packet serialization through LeaderElection GenServer
|
||||||
- [ ] Add backpressure mechanism — no global rate limiting if APRS-IS sends burst traffic; only defense is fixed-size buffer with silent drops
|
- [ ] Add backpressure mechanism — no global rate limiting if APRS-IS sends burst traffic; only defense is fixed-size buffer with silent drops
|
||||||
- [ ] Cluster packet distribution race — `PacketDistributor.distribute_packet/1` only broadcasts if currently leader; leadership change between receipt and broadcast drops packets
|
|
||||||
|
|
||||||
## Front-End Display
|
## Front-End Display
|
||||||
|
|
||||||
- [x] Fix XSS vulnerability in PopupComponent fallback — added `escapeHtml()` to `map_helpers.ts`, applied to callsign and comment in `buildPopupContent` in `map.ts`
|
- [x] Fix XSS vulnerability in PopupComponent fallback — added `escapeHtml()` to `map_helpers.ts`, applied to callsign and comment in `buildPopupContent` in `map.ts`
|
||||||
- [ ] Simplify coordinate extraction in `packets_live/index.html.heex:65-134` — deeply nested conditional logic with multiple fallback chains
|
- [x] Simplify coordinate extraction in `packets_live/index.html.heex` — extracted `extract_coordinate/2` and `format_coordinate/1` helpers into `PacketsLive.Index`; template went from 70 lines of nested conditionals to 6 lines
|
||||||
- [x] Fix memory leak in InfoMap hook — stored `setTimeout` ref in `this.resizeTimer`, cancel in `destroyed()`
|
- [x] Fix memory leak in InfoMap hook — stored `setTimeout` ref in `this.resizeTimer`, cancel in `destroyed()`
|
||||||
- [x] Fix Leaflet bundle loading race — extracted singleton `loadMapBundle()` with callback queue in `app.js`
|
- [x] Fix Leaflet bundle loading race — extracted singleton `loadMapBundle()` with callback queue in `app.js`
|
||||||
- [ ] Add loading indicator for real-time bounds updates — only `@historical_loading` triggers spinner, not bounds filtering
|
- [ ] Add loading indicator for real-time bounds updates — only `@historical_loading` triggers spinner, not bounds filtering
|
||||||
- [x] Fix stale generation check bypass in `historical_loader.ex:100-108` — split into two function clauses: nil generation always loads, integer generation checks staleness
|
- [x] Fix stale generation check bypass in `historical_loader.ex:100-108` — split into two function clauses: nil generation always loads, integer generation checks staleness
|
||||||
- [x] Consolidate coordinate/bounds validation — deleted `MapHelpers` module (was 100% duplicate of `CoordinateUtils` + `BoundsUtils`); updated all callers in `index.ex`, `data_builder.ex`, `mobile_channel.ex`
|
- [x] Consolidate coordinate/bounds validation — deleted `MapHelpers` module (was 100% duplicate of `CoordinateUtils` + `BoundsUtils`); updated all callers in `index.ex`, `data_builder.ex`, `mobile_channel.ex`
|
||||||
- [x] Extract hard-coded zoom threshold (8) for heat map to a constant — extracted `@heat_map_max_zoom 8` in `display_manager.ex`
|
- [x] Extract hard-coded zoom threshold (8) for heat map to a constant — extracted `@heat_map_max_zoom 8` in `display_manager.ex`
|
||||||
|
- [x] Remove debug logging from hot paths — removed `Logger.debug` calls from `packets.ex` query path, `PacketDistributor`, and `console.log` statements from `app.js`
|
||||||
|
|
||||||
## Packet Purging
|
## Packet Purging
|
||||||
|
|
||||||
- [x] Shorter default retention — changed from 365 days to 7 days (configurable via `PACKET_RETENTION_DAYS` env var)
|
- [x] Shorter default retention — changed from 365 days to 7 days (configurable via `PACKET_RETENTION_DAYS` env var)
|
||||||
- [x] Improve cleanup efficiency — replaced two-step SELECT IDs + DELETE by IDs with single-query CTE-based batch DELETE; eliminates extra round-trip per batch
|
- [x] Improve cleanup efficiency — replaced two-step SELECT IDs + DELETE by IDs with single-query CTE-based batch DELETE; eliminates extra round-trip per batch
|
||||||
- [x] Add cleanup telemetry — added `:telemetry.execute` to `cleanup_packets_older_than_batched/1`
|
- [x] Add cleanup telemetry — added `:telemetry.execute` to `cleanup_packets_older_than_batched/1`
|
||||||
- [ ] Add partial index for cleanup queries — `WHERE received_at < cutoff` would benefit from a partial index on old packets
|
- [~] Partial index for cleanup queries — assessed: existing `packets_received_at_idx` B-tree is already optimal for `WHERE received_at < $1 LIMIT $2`; partial index with dynamic cutoff adds no benefit
|
||||||
- [~] ETS PacketStore TTL mismatch — assessed: the 2-hour ETS TTL is intentional for LiveView memory efficiency; DB retention is for historical data. Different purposes, not a bug.
|
- [~] ETS PacketStore TTL mismatch — assessed: the 2-hour ETS TTL is intentional for LiveView memory efficiency; DB retention is for historical data. Different purposes, not a bug.
|
||||||
- [ ] Consider PostgreSQL table partitioning — partition packets by time range (daily/weekly) for instant `DROP PARTITION` cleanup instead of batch DELETEs; requires one-time migration
|
- [ ] Consider PostgreSQL table partitioning — partition packets by time range (daily/weekly) for instant `DROP PARTITION` cleanup instead of batch DELETEs; requires one-time migration but not worth the complexity unless cleanup exceeds 5-minute time limit regularly
|
||||||
|
|
|
||||||
|
|
@ -1,5 +1,3 @@
|
||||||
console.log("app.js loading...");
|
|
||||||
|
|
||||||
// If you want to use Phoenix channels, run `mix help phx.gen.channel`
|
// If you want to use Phoenix channels, run `mix help phx.gen.channel`
|
||||||
// to get started and then uncomment the line below.
|
// to get started and then uncomment the line below.
|
||||||
// import "./user_socket.js"
|
// import "./user_socket.js"
|
||||||
|
|
@ -86,10 +84,8 @@ const originalMapMounted = MapAPRSMap.mounted;
|
||||||
Hooks.APRSMap = {
|
Hooks.APRSMap = {
|
||||||
...MapAPRSMap,
|
...MapAPRSMap,
|
||||||
mounted() {
|
mounted() {
|
||||||
console.log("APRSMap wrapper mounted() called");
|
|
||||||
const self = this;
|
const self = this;
|
||||||
loadMapBundle(() => {
|
loadMapBundle(() => {
|
||||||
console.log("Map bundle ready, calling original mounted");
|
|
||||||
if (originalMapMounted) {
|
if (originalMapMounted) {
|
||||||
originalMapMounted.call(self);
|
originalMapMounted.call(self);
|
||||||
}
|
}
|
||||||
|
|
@ -180,7 +176,6 @@ window.matchMedia("(prefers-color-scheme: dark)").addEventListener("change", ()
|
||||||
}
|
}
|
||||||
});
|
});
|
||||||
|
|
||||||
console.log("Creating LiveSocket with hooks:", Object.keys(Hooks));
|
|
||||||
let liveSocket = new LiveSocket("/live", Socket, {
|
let liveSocket = new LiveSocket("/live", Socket, {
|
||||||
longPollFallbackMs: 2500,
|
longPollFallbackMs: 2500,
|
||||||
params: { _csrf_token: csrfToken, viewport_width: window.innerWidth },
|
params: { _csrf_token: csrfToken, viewport_width: window.innerWidth },
|
||||||
|
|
@ -196,7 +191,6 @@ window.addEventListener("phx:page-loading-stop", (_info) => topbar.hide());
|
||||||
// Handle connection draining reconnect events
|
// Handle connection draining reconnect events
|
||||||
window.addEventListener("phx:reconnect", (e) => {
|
window.addEventListener("phx:reconnect", (e) => {
|
||||||
const delay = e.detail.delay || 1000;
|
const delay = e.detail.delay || 1000;
|
||||||
console.log(`[LiveSocket] Reconnecting in ${delay}ms due to connection draining...`);
|
|
||||||
setTimeout(() => {
|
setTimeout(() => {
|
||||||
// Disconnect and reconnect to potentially land on a different server
|
// Disconnect and reconnect to potentially land on a different server
|
||||||
liveSocket.disconnect();
|
liveSocket.disconnect();
|
||||||
|
|
@ -216,7 +210,6 @@ window.addEventListener("phx:live_socket:connect", (info) => {
|
||||||
if (socket && socket.fallbackTimer) {
|
if (socket && socket.fallbackTimer) {
|
||||||
clearTimeout(socket.fallbackTimer);
|
clearTimeout(socket.fallbackTimer);
|
||||||
socket.fallbackTimer = null;
|
socket.fallbackTimer = null;
|
||||||
console.log("[LiveSocket] Cleared fallback timer after successful connection");
|
|
||||||
}
|
}
|
||||||
});
|
});
|
||||||
|
|
||||||
|
|
@ -226,7 +219,6 @@ setTimeout(() => {
|
||||||
if (socket && socket.isConnected() && socket.fallbackTimer) {
|
if (socket && socket.isConnected() && socket.fallbackTimer) {
|
||||||
clearTimeout(socket.fallbackTimer);
|
clearTimeout(socket.fallbackTimer);
|
||||||
socket.fallbackTimer = null;
|
socket.fallbackTimer = null;
|
||||||
console.log("[LiveSocket] Cleared lingering fallback timer");
|
|
||||||
}
|
}
|
||||||
}, 5000);
|
}, 5000);
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -20,6 +20,15 @@ defmodule Aprsme.Cluster.LeaderElection do
|
||||||
GenServer.call(__MODULE__, :leader?)
|
GenServer.call(__MODULE__, :leader?)
|
||||||
end
|
end
|
||||||
|
|
||||||
|
@doc """
|
||||||
|
Fast cached leadership check using :persistent_term.
|
||||||
|
No GenServer.call overhead — suitable for hot paths like packet distribution.
|
||||||
|
"""
|
||||||
|
@spec leader_cached?() :: boolean()
|
||||||
|
def leader_cached? do
|
||||||
|
:persistent_term.get({__MODULE__, :is_leader}, false)
|
||||||
|
end
|
||||||
|
|
||||||
def current_leader do
|
def current_leader do
|
||||||
GenServer.call(__MODULE__, :current_leader)
|
GenServer.call(__MODULE__, :current_leader)
|
||||||
end
|
end
|
||||||
|
|
@ -60,6 +69,9 @@ defmodule Aprsme.Cluster.LeaderElection do
|
||||||
# Schedule periodic checks
|
# Schedule periodic checks
|
||||||
Process.send_after(self(), :check_leadership, @check_interval)
|
Process.send_after(self(), :check_leadership, @check_interval)
|
||||||
|
|
||||||
|
# Initialize cached leadership state
|
||||||
|
:persistent_term.put({__MODULE__, :is_leader}, false)
|
||||||
|
|
||||||
{:ok, %{is_leader: false, leader_node: nil, cluster_enabled: cluster_enabled, election_forced: false}}
|
{:ok, %{is_leader: false, leader_node: nil, cluster_enabled: cluster_enabled, election_forced: false}}
|
||||||
end
|
end
|
||||||
|
|
||||||
|
|
@ -117,6 +129,7 @@ defmodule Aprsme.Cluster.LeaderElection do
|
||||||
case :global.register_name(@election_key, self(), &resolve_conflict/3) do
|
case :global.register_name(@election_key, self(), &resolve_conflict/3) do
|
||||||
:yes ->
|
:yes ->
|
||||||
Logger.info("Elected as APRS-IS connection leader on node #{node()}")
|
Logger.info("Elected as APRS-IS connection leader on node #{node()}")
|
||||||
|
:persistent_term.put({__MODULE__, :is_leader}, true)
|
||||||
notify_leadership_change(true)
|
notify_leadership_change(true)
|
||||||
{:noreply, %{state | is_leader: true, leader_node: node()}}
|
{:noreply, %{state | is_leader: true, leader_node: node()}}
|
||||||
|
|
||||||
|
|
@ -124,6 +137,7 @@ defmodule Aprsme.Cluster.LeaderElection do
|
||||||
leader_pid = :global.whereis_name(@election_key)
|
leader_pid = :global.whereis_name(@election_key)
|
||||||
leader_node = if leader_pid != :undefined and is_pid(leader_pid), do: node(leader_pid)
|
leader_node = if leader_pid != :undefined and is_pid(leader_pid), do: node(leader_pid)
|
||||||
Logger.info("Not elected as leader. Current leader is on node #{inspect(leader_node)}")
|
Logger.info("Not elected as leader. Current leader is on node #{inspect(leader_node)}")
|
||||||
|
:persistent_term.put({__MODULE__, :is_leader}, false)
|
||||||
{:noreply, %{state | is_leader: false, leader_node: leader_node}}
|
{:noreply, %{state | is_leader: false, leader_node: leader_node}}
|
||||||
end
|
end
|
||||||
end
|
end
|
||||||
|
|
@ -161,6 +175,7 @@ defmodule Aprsme.Cluster.LeaderElection do
|
||||||
def terminate(reason, state) do
|
def terminate(reason, state) do
|
||||||
if state.is_leader do
|
if state.is_leader do
|
||||||
Logger.info("Leader stepping down due to: #{inspect(reason)}")
|
Logger.info("Leader stepping down due to: #{inspect(reason)}")
|
||||||
|
:persistent_term.put({__MODULE__, :is_leader}, false)
|
||||||
:global.unregister_name(@election_key)
|
:global.unregister_name(@election_key)
|
||||||
notify_leadership_change(false)
|
notify_leadership_change(false)
|
||||||
end
|
end
|
||||||
|
|
|
||||||
|
|
@ -7,23 +7,19 @@ defmodule Aprsme.Cluster.PacketDistributor do
|
||||||
alias Aprsme.Cluster.LeaderElection
|
alias Aprsme.Cluster.LeaderElection
|
||||||
alias AprsmeWeb.MapLive.PacketStore
|
alias AprsmeWeb.MapLive.PacketStore
|
||||||
|
|
||||||
require Logger
|
|
||||||
|
|
||||||
@pubsub_topic "cluster:packets"
|
@pubsub_topic "cluster:packets"
|
||||||
|
|
||||||
def distribute_packet(packet) do
|
def distribute_packet(packet) do
|
||||||
# Only distribute if clustering is enabled and we're the leader
|
# Only distribute if clustering is enabled and we're the leader
|
||||||
cluster_enabled = Application.get_env(:aprsme, :cluster_enabled, false)
|
cluster_enabled = Application.get_env(:aprsme, :cluster_enabled, false)
|
||||||
|
|
||||||
if cluster_enabled and LeaderElection.leader?() do
|
if cluster_enabled and LeaderElection.leader_cached?() do
|
||||||
# Broadcast to all nodes including self
|
# Broadcast to all nodes including self
|
||||||
Phoenix.PubSub.broadcast(
|
Phoenix.PubSub.broadcast(
|
||||||
Aprsme.PubSub,
|
Aprsme.PubSub,
|
||||||
@pubsub_topic,
|
@pubsub_topic,
|
||||||
{:distributed_packet, packet}
|
{:distributed_packet, packet}
|
||||||
)
|
)
|
||||||
|
|
||||||
Logger.debug("Distributed packet #{packet.raw} to cluster")
|
|
||||||
end
|
end
|
||||||
end
|
end
|
||||||
|
|
||||||
|
|
@ -38,6 +34,6 @@ defmodule Aprsme.Cluster.PacketDistributor do
|
||||||
# Update packet store for LiveView
|
# Update packet store for LiveView
|
||||||
PacketStore.store_packet(packet)
|
PacketStore.store_packet(packet)
|
||||||
|
|
||||||
Logger.debug("Received distributed packet on node #{node()}")
|
:ok
|
||||||
end
|
end
|
||||||
end
|
end
|
||||||
|
|
|
||||||
|
|
@ -2,6 +2,9 @@ defmodule Aprsme.PacketProducer do
|
||||||
@moduledoc """
|
@moduledoc """
|
||||||
GenStage producer that handles incoming APRS packets and sends them to consumers
|
GenStage producer that handles incoming APRS packets and sends them to consumers
|
||||||
for efficient batch processing.
|
for efficient batch processing.
|
||||||
|
|
||||||
|
Uses `:queue` for O(1) enqueue/dequeue instead of lists, which avoids
|
||||||
|
O(n) `Enum.take` on every buffer overflow during traffic bursts.
|
||||||
"""
|
"""
|
||||||
use GenStage
|
use GenStage
|
||||||
|
|
||||||
|
|
@ -17,14 +20,15 @@ defmodule Aprsme.PacketProducer do
|
||||||
|
|
||||||
@impl true
|
@impl true
|
||||||
def init(opts) do
|
def init(opts) do
|
||||||
{:producer, %{demand: 0, buffer: [], buffer_size: 0, max_buffer_size: opts[:max_buffer_size] || 1000}}
|
{:producer, %{demand: 0, buffer: :queue.new(), buffer_size: 0, max_buffer_size: opts[:max_buffer_size] || 1000}}
|
||||||
end
|
end
|
||||||
|
|
||||||
@impl true
|
@impl true
|
||||||
def handle_demand(incoming_demand, %{demand: demand, buffer: buffer} = state) do
|
def handle_demand(incoming_demand, %{demand: demand} = state) do
|
||||||
{events, remaining_buffer, remaining_demand} = dispatch_events(buffer, demand + incoming_demand)
|
{events, new_buffer, remaining_demand} = dispatch_events(state.buffer, demand + incoming_demand)
|
||||||
remaining_size = state.buffer_size - length(events)
|
|
||||||
{:noreply, events, %{state | demand: remaining_demand, buffer: remaining_buffer, buffer_size: remaining_size}}
|
{:noreply, events,
|
||||||
|
%{state | demand: remaining_demand, buffer: new_buffer, buffer_size: state.buffer_size - length(events)}}
|
||||||
end
|
end
|
||||||
|
|
||||||
@impl true
|
@impl true
|
||||||
|
|
@ -51,16 +55,23 @@ defmodule Aprsme.PacketProducer do
|
||||||
%{}
|
%{}
|
||||||
)
|
)
|
||||||
|
|
||||||
{:noreply, [], %{state | buffer: [packet_data | Enum.take(buffer, max_size - 1)], buffer_size: max_size}}
|
# Drop oldest (front of queue), add new to back — O(1) amortized
|
||||||
|
{_, trimmed} = :queue.out(buffer)
|
||||||
|
{:noreply, [], %{state | buffer: :queue.in(packet_data, trimmed), buffer_size: max_size}}
|
||||||
else
|
else
|
||||||
{:noreply, [], %{state | buffer: [packet_data | buffer], buffer_size: new_size}}
|
{:noreply, [], %{state | buffer: :queue.in(packet_data, buffer), buffer_size: new_size}}
|
||||||
end
|
end
|
||||||
end
|
end
|
||||||
end
|
end
|
||||||
|
|
||||||
defp dispatch_events(buffer, demand) when demand > 0 and buffer != [] do
|
defp dispatch_events(buffer, demand) when demand > 0 do
|
||||||
{events, remaining} = Enum.split(buffer, demand)
|
if :queue.is_empty(buffer) do
|
||||||
{events, remaining, demand - length(events)}
|
{[], buffer, demand}
|
||||||
|
else
|
||||||
|
{front, remaining} = :queue.split(min(demand, :queue.len(buffer)), buffer)
|
||||||
|
events = :queue.to_list(front)
|
||||||
|
{events, remaining, demand - length(events)}
|
||||||
|
end
|
||||||
end
|
end
|
||||||
|
|
||||||
defp dispatch_events(buffer, demand) do
|
defp dispatch_events(buffer, demand) do
|
||||||
|
|
|
||||||
|
|
@ -225,13 +225,6 @@ defmodule Aprsme.Packets do
|
||||||
# Ensure data_extended is properly sanitized before insertion
|
# Ensure data_extended is properly sanitized before insertion
|
||||||
attrs = sanitize_data_extended_attr(attrs)
|
attrs = sanitize_data_extended_attr(attrs)
|
||||||
|
|
||||||
# Debug log to see what we're trying to insert
|
|
||||||
if attrs[:data_extended] do
|
|
||||||
require Logger
|
|
||||||
|
|
||||||
Logger.debug("Final data_extended before insert: #{inspect(attrs[:data_extended], binaries: :as_binaries)}")
|
|
||||||
end
|
|
||||||
|
|
||||||
case %Packet{} |> Packet.changeset(attrs) |> Repo.insert() do
|
case %Packet{} |> Packet.changeset(attrs) |> Repo.insert() do
|
||||||
{:ok, packet} ->
|
{:ok, packet} ->
|
||||||
# Invalidate cache for this packet's callsign
|
# Invalidate cache for this packet's callsign
|
||||||
|
|
@ -502,8 +495,6 @@ defmodule Aprsme.Packets do
|
||||||
def get_recent_packets(opts \\ %{}) do
|
def get_recent_packets(opts \\ %{}) do
|
||||||
require Logger
|
require Logger
|
||||||
|
|
||||||
Logger.debug("Packets.get_recent_packets called with opts: #{inspect(opts)}")
|
|
||||||
|
|
||||||
# Use hours_back from opts if provided, otherwise default to 24 hours
|
# Use hours_back from opts if provided, otherwise default to 24 hours
|
||||||
hours_back = Map.get(opts, :hours_back, 24)
|
hours_back = Map.get(opts, :hours_back, 24)
|
||||||
time_ago = DateTime.add(DateTime.utc_now(), -hours_back * 3600, :second)
|
time_ago = DateTime.add(DateTime.utc_now(), -hours_back * 3600, :second)
|
||||||
|
|
@ -544,9 +535,7 @@ defmodule Aprsme.Packets do
|
||||||
|> offset(^offset)
|
|> offset(^offset)
|
||||||
|> QueryBuilder.with_coordinates()
|
|> QueryBuilder.with_coordinates()
|
||||||
|
|
||||||
result = Repo.all(query)
|
Repo.all(query)
|
||||||
Logger.debug("Packets.get_recent_packets returning #{length(result)} packets")
|
|
||||||
result
|
|
||||||
end
|
end
|
||||||
|
|
||||||
@doc """
|
@doc """
|
||||||
|
|
|
||||||
|
|
@ -32,4 +32,45 @@ defmodule AprsmeWeb.PacketsLive.Index do
|
||||||
socket = assign(socket, :packets, packets)
|
socket = assign(socket, :packets, packets)
|
||||||
{:noreply, socket}
|
{:noreply, socket}
|
||||||
end
|
end
|
||||||
|
|
||||||
|
@doc """
|
||||||
|
Extract a coordinate (:lat or :lon) from a packet, checking multiple sources:
|
||||||
|
1. Direct :lat/:lon keys
|
||||||
|
2. Packet struct with PostGIS location
|
||||||
|
3. data_extended map with :latitude/:longitude keys
|
||||||
|
"""
|
||||||
|
def extract_coordinate(packet, which) when which in [:lat, :lon] do
|
||||||
|
direct_key = which
|
||||||
|
extended_key = if which == :lat, do: :latitude, else: :longitude
|
||||||
|
|
||||||
|
Map.get(packet, direct_key) ||
|
||||||
|
extract_from_location(packet, which) ||
|
||||||
|
extract_from_data_extended(packet, extended_key)
|
||||||
|
end
|
||||||
|
|
||||||
|
@doc """
|
||||||
|
Format a coordinate value for display with up to 6 decimal places.
|
||||||
|
"""
|
||||||
|
def format_coordinate(nil), do: ""
|
||||||
|
|
||||||
|
def format_coordinate(value) when is_float(value) do
|
||||||
|
"~.6f" |> :io_lib.format([value]) |> List.to_string()
|
||||||
|
end
|
||||||
|
|
||||||
|
def format_coordinate(value) when is_binary(value) do
|
||||||
|
Regex.replace(~r/(\d+\.\d{1,6})\d*/, value, "\\1")
|
||||||
|
end
|
||||||
|
|
||||||
|
def format_coordinate(value), do: to_string(value)
|
||||||
|
|
||||||
|
defp extract_from_location(%Aprsme.Packet{location: %Geo.Point{coordinates: {_lon, lat}}}, :lat), do: lat
|
||||||
|
defp extract_from_location(%Aprsme.Packet{location: %Geo.Point{coordinates: {lon, _lat}}}, :lon), do: lon
|
||||||
|
defp extract_from_location(_, _), do: nil
|
||||||
|
|
||||||
|
defp extract_from_data_extended(packet, key) do
|
||||||
|
case Map.get(packet, :data_extended) do
|
||||||
|
%{} = data -> Map.get(data, key) || Map.get(data, to_string(key))
|
||||||
|
_ -> nil
|
||||||
|
end
|
||||||
|
end
|
||||||
end
|
end
|
||||||
|
|
|
||||||
|
|
@ -64,72 +64,12 @@
|
||||||
</td>
|
</td>
|
||||||
<td>
|
<td>
|
||||||
<span class="text-xs font-mono">
|
<span class="text-xs font-mono">
|
||||||
<% lat =
|
{format_coordinate(extract_coordinate(packet, :lat))}
|
||||||
Map.get(packet, :lat) ||
|
|
||||||
if(
|
|
||||||
Map.has_key?(packet, :location) and not is_nil(packet.location) and
|
|
||||||
function_exported?(Aprsme.Packet, :lat, 1),
|
|
||||||
do: Aprsme.Packet.lat(packet),
|
|
||||||
else: nil
|
|
||||||
) ||
|
|
||||||
case Map.get(packet, :data_extended) do
|
|
||||||
%{} = data_ext ->
|
|
||||||
case Map.get(data_ext, :latitude) || Map.get(data_ext, "latitude") do
|
|
||||||
{:ok, value} -> value
|
|
||||||
value -> value
|
|
||||||
end
|
|
||||||
|
|
||||||
_ ->
|
|
||||||
nil
|
|
||||||
end %>
|
|
||||||
<%= if not is_nil(lat) do %>
|
|
||||||
<%= if is_float(lat) do %>
|
|
||||||
{:io_lib.format("~.6f", [lat]) |> List.to_string()}
|
|
||||||
<% else %>
|
|
||||||
<%= if is_binary(lat) do %>
|
|
||||||
{Regex.replace(~r/(\d+\.\d{1,6})\d*/, lat, "\\1")}
|
|
||||||
<% else %>
|
|
||||||
{lat}
|
|
||||||
<% end %>
|
|
||||||
<% end %>
|
|
||||||
<% else %>
|
|
||||||
""
|
|
||||||
<% end %>
|
|
||||||
</span>
|
</span>
|
||||||
</td>
|
</td>
|
||||||
<td>
|
<td>
|
||||||
<span class="text-xs font-mono">
|
<span class="text-xs font-mono">
|
||||||
<% lon =
|
{format_coordinate(extract_coordinate(packet, :lon))}
|
||||||
Map.get(packet, :lon) ||
|
|
||||||
if(
|
|
||||||
Map.has_key?(packet, :location) and not is_nil(packet.location) and
|
|
||||||
function_exported?(Aprsme.Packet, :lon, 1),
|
|
||||||
do: Aprsme.Packet.lon(packet),
|
|
||||||
else: nil
|
|
||||||
) ||
|
|
||||||
case Map.get(packet, :data_extended) do
|
|
||||||
%{} = data_ext ->
|
|
||||||
case Map.get(data_ext, :longitude) || Map.get(data_ext, "longitude") do
|
|
||||||
{:ok, value} -> value
|
|
||||||
value -> value
|
|
||||||
end
|
|
||||||
|
|
||||||
_ ->
|
|
||||||
nil
|
|
||||||
end %>
|
|
||||||
<%= if not is_nil(lon) do %>
|
|
||||||
<%= if is_float(lon) do %>
|
|
||||||
{:io_lib.format("~.6f", [lon]) |> List.to_string()}
|
|
||||||
<% else %>
|
|
||||||
<%= if is_binary(lon) do %>
|
|
||||||
{Regex.replace(~r/(\d+\.\d{1,6})\d*/, lon, "\\1")}
|
|
||||||
<% else %>
|
|
||||||
{lon}
|
|
||||||
<% end %>
|
|
||||||
<% end %>
|
|
||||||
<% else %>
|
|
||||||
""
|
|
||||||
<% end %>
|
|
||||||
</span>
|
</span>
|
||||||
</td>
|
</td>
|
||||||
<td>
|
<td>
|
||||||
|
|
|
||||||
1
mix.exs
1
mix.exs
|
|
@ -96,7 +96,6 @@ defmodule Aprsme.MixProject do
|
||||||
{:floki, ">= 0.30.0", only: :test},
|
{:floki, ">= 0.30.0", only: :test},
|
||||||
{:lazy_html, ">= 0.1.0", only: :test},
|
{:lazy_html, ">= 0.1.0", only: :test},
|
||||||
{:mix_test_watch, "~> 1.1", only: [:dev, :test]},
|
{:mix_test_watch, "~> 1.1", only: [:dev, :test]},
|
||||||
{:sobelow, "~> 0.8", only: :dev},
|
|
||||||
{:stream_data, "~> 1.2.0", only: [:dev, :test]},
|
{:stream_data, "~> 1.2.0", only: [:dev, :test]},
|
||||||
{:igniter, "~> 0.7.0", only: [:dev, :test]},
|
{:igniter, "~> 0.7.0", only: [:dev, :test]},
|
||||||
{:mox, "~> 1.2", only: :test},
|
{:mox, "~> 1.2", only: :test},
|
||||||
|
|
|
||||||
|
|
@ -78,6 +78,22 @@ defmodule Aprsme.Cluster.LeaderElectionTest do
|
||||||
test "returns leadership status" do
|
test "returns leadership status" do
|
||||||
assert is_boolean(LeaderElection.leader?())
|
assert is_boolean(LeaderElection.leader?())
|
||||||
end
|
end
|
||||||
|
|
||||||
|
test "cached check matches GenServer state" do
|
||||||
|
# leader_cached?/0 reads from :persistent_term — no GenServer.call
|
||||||
|
assert LeaderElection.leader_cached?() == LeaderElection.leader?()
|
||||||
|
end
|
||||||
|
|
||||||
|
test "cached check returns false before election" do
|
||||||
|
# Stop current instance
|
||||||
|
GenServer.stop(LeaderElection)
|
||||||
|
:global.unregister_name({:aprs_is_leader, LeaderElection})
|
||||||
|
|
||||||
|
# Clear the persistent_term
|
||||||
|
:persistent_term.put({LeaderElection, :is_leader}, false)
|
||||||
|
|
||||||
|
assert LeaderElection.leader_cached?() == false
|
||||||
|
end
|
||||||
end
|
end
|
||||||
|
|
||||||
describe "current_leader/0" do
|
describe "current_leader/0" do
|
||||||
|
|
|
||||||
|
|
@ -84,7 +84,7 @@ defmodule Aprsme.Cluster.PacketDistributorTest do
|
||||||
# StreamingPacketsPubSub and PacketStore are already running in test
|
# StreamingPacketsPubSub and PacketStore are already running in test
|
||||||
result = PacketDistributor.handle_distributed_packet({:distributed_packet, @test_packet})
|
result = PacketDistributor.handle_distributed_packet({:distributed_packet, @test_packet})
|
||||||
|
|
||||||
# The function logs and returns :ok from Logger.debug
|
# Broadcasts to local clients and stores in PacketStore
|
||||||
assert result == :ok
|
assert result == :ok
|
||||||
end
|
end
|
||||||
end
|
end
|
||||||
|
|
|
||||||
112
test/aprsme/packet_producer_test.exs
Normal file
112
test/aprsme/packet_producer_test.exs
Normal file
|
|
@ -0,0 +1,112 @@
|
||||||
|
defmodule Aprsme.PacketProducerTest do
|
||||||
|
use ExUnit.Case, async: true
|
||||||
|
|
||||||
|
alias Aprsme.PacketProducer
|
||||||
|
|
||||||
|
describe "init/1" do
|
||||||
|
test "initializes with default max_buffer_size" do
|
||||||
|
{:producer, state} = PacketProducer.init([])
|
||||||
|
assert state.demand == 0
|
||||||
|
assert state.buffer_size == 0
|
||||||
|
assert state.max_buffer_size == 1000
|
||||||
|
end
|
||||||
|
|
||||||
|
test "initializes with custom max_buffer_size" do
|
||||||
|
{:producer, state} = PacketProducer.init(max_buffer_size: 500)
|
||||||
|
assert state.max_buffer_size == 500
|
||||||
|
end
|
||||||
|
|
||||||
|
test "initializes with empty :queue buffer" do
|
||||||
|
{:producer, state} = PacketProducer.init([])
|
||||||
|
assert :queue.is_queue(state.buffer)
|
||||||
|
assert :queue.is_empty(state.buffer)
|
||||||
|
end
|
||||||
|
end
|
||||||
|
|
||||||
|
describe "handle_cast {:packet, _} with demand > 0" do
|
||||||
|
test "dispatches packet immediately when demand exists" do
|
||||||
|
state = %{demand: 3, buffer: :queue.new(), buffer_size: 0, max_buffer_size: 1000}
|
||||||
|
{:noreply, events, new_state} = PacketProducer.handle_cast({:packet, %{sender: "TEST"}}, state)
|
||||||
|
|
||||||
|
assert events == [%{sender: "TEST"}]
|
||||||
|
assert new_state.demand == 2
|
||||||
|
end
|
||||||
|
end
|
||||||
|
|
||||||
|
describe "handle_cast {:packet, _} with demand == 0" do
|
||||||
|
test "buffers packet when no demand" do
|
||||||
|
state = %{demand: 0, buffer: :queue.new(), buffer_size: 0, max_buffer_size: 1000}
|
||||||
|
{:noreply, [], new_state} = PacketProducer.handle_cast({:packet, %{sender: "TEST"}}, state)
|
||||||
|
|
||||||
|
assert new_state.buffer_size == 1
|
||||||
|
assert :queue.len(new_state.buffer) == 1
|
||||||
|
end
|
||||||
|
|
||||||
|
test "drops oldest packet when buffer overflows" do
|
||||||
|
# Fill buffer to max
|
||||||
|
buffer =
|
||||||
|
Enum.reduce(1..3, :queue.new(), fn i, q ->
|
||||||
|
:queue.in(%{sender: "PACKET#{i}"}, q)
|
||||||
|
end)
|
||||||
|
|
||||||
|
state = %{demand: 0, buffer: buffer, buffer_size: 3, max_buffer_size: 3}
|
||||||
|
|
||||||
|
{:noreply, [], new_state} =
|
||||||
|
PacketProducer.handle_cast({:packet, %{sender: "NEW"}}, state)
|
||||||
|
|
||||||
|
# Buffer size should still be max
|
||||||
|
assert new_state.buffer_size == 3
|
||||||
|
|
||||||
|
# Oldest packet (PACKET1) should be dropped, NEW should be present
|
||||||
|
items = :queue.to_list(new_state.buffer)
|
||||||
|
senders = Enum.map(items, & &1.sender)
|
||||||
|
refute "PACKET1" in senders
|
||||||
|
assert "NEW" in senders
|
||||||
|
assert "PACKET2" in senders
|
||||||
|
assert "PACKET3" in senders
|
||||||
|
end
|
||||||
|
end
|
||||||
|
|
||||||
|
describe "handle_demand/2" do
|
||||||
|
test "dispatches buffered packets in FIFO order" do
|
||||||
|
buffer =
|
||||||
|
Enum.reduce(1..5, :queue.new(), fn i, q ->
|
||||||
|
:queue.in(%{sender: "P#{i}"}, q)
|
||||||
|
end)
|
||||||
|
|
||||||
|
state = %{demand: 0, buffer: buffer, buffer_size: 5, max_buffer_size: 1000}
|
||||||
|
{:noreply, events, new_state} = PacketProducer.handle_demand(3, state)
|
||||||
|
|
||||||
|
# Should dispatch first 3 (FIFO order)
|
||||||
|
assert length(events) == 3
|
||||||
|
assert Enum.map(events, & &1.sender) == ["P1", "P2", "P3"]
|
||||||
|
assert new_state.buffer_size == 2
|
||||||
|
assert new_state.demand == 0
|
||||||
|
end
|
||||||
|
|
||||||
|
test "dispatches all available when demand exceeds buffer" do
|
||||||
|
buffer = :queue.in(%{sender: "ONLY"}, :queue.new())
|
||||||
|
|
||||||
|
state = %{demand: 0, buffer: buffer, buffer_size: 1, max_buffer_size: 1000}
|
||||||
|
{:noreply, events, new_state} = PacketProducer.handle_demand(5, state)
|
||||||
|
|
||||||
|
assert events == [%{sender: "ONLY"}]
|
||||||
|
assert new_state.buffer_size == 0
|
||||||
|
assert new_state.demand == 4
|
||||||
|
end
|
||||||
|
|
||||||
|
test "stores demand when buffer is empty" do
|
||||||
|
state = %{demand: 0, buffer: :queue.new(), buffer_size: 0, max_buffer_size: 1000}
|
||||||
|
{:noreply, [], new_state} = PacketProducer.handle_demand(10, state)
|
||||||
|
|
||||||
|
assert new_state.demand == 10
|
||||||
|
end
|
||||||
|
|
||||||
|
test "accumulates demand" do
|
||||||
|
state = %{demand: 5, buffer: :queue.new(), buffer_size: 0, max_buffer_size: 1000}
|
||||||
|
{:noreply, [], new_state} = PacketProducer.handle_demand(3, state)
|
||||||
|
|
||||||
|
assert new_state.demand == 8
|
||||||
|
end
|
||||||
|
end
|
||||||
|
end
|
||||||
70
test/aprsme_web/live/packets_live/index_test.exs
Normal file
70
test/aprsme_web/live/packets_live/index_test.exs
Normal file
|
|
@ -0,0 +1,70 @@
|
||||||
|
defmodule AprsmeWeb.PacketsLive.IndexTest do
|
||||||
|
use ExUnit.Case, async: true
|
||||||
|
|
||||||
|
alias AprsmeWeb.PacketsLive.Index
|
||||||
|
|
||||||
|
describe "extract_coordinate/2" do
|
||||||
|
test "extracts lat from direct :lat key" do
|
||||||
|
packet = %{lat: 35.123456, lon: -75.654321}
|
||||||
|
assert Index.extract_coordinate(packet, :lat) == 35.123456
|
||||||
|
end
|
||||||
|
|
||||||
|
test "extracts lon from direct :lon key" do
|
||||||
|
packet = %{lat: 35.0, lon: -75.654321}
|
||||||
|
assert Index.extract_coordinate(packet, :lon) == -75.654321
|
||||||
|
end
|
||||||
|
|
||||||
|
test "extracts lat from Packet struct with location" do
|
||||||
|
packet = %Aprsme.Packet{location: %Geo.Point{coordinates: {-75.0, 35.0}, srid: 4326}}
|
||||||
|
assert Index.extract_coordinate(packet, :lat) == 35.0
|
||||||
|
end
|
||||||
|
|
||||||
|
test "extracts lon from Packet struct with location" do
|
||||||
|
packet = %Aprsme.Packet{location: %Geo.Point{coordinates: {-75.0, 35.0}, srid: 4326}}
|
||||||
|
assert Index.extract_coordinate(packet, :lon) == -75.0
|
||||||
|
end
|
||||||
|
|
||||||
|
test "extracts lat from data_extended latitude" do
|
||||||
|
packet = %{data_extended: %{latitude: 35.5, longitude: -75.5}}
|
||||||
|
assert Index.extract_coordinate(packet, :lat) == 35.5
|
||||||
|
end
|
||||||
|
|
||||||
|
test "extracts lon from data_extended longitude" do
|
||||||
|
packet = %{data_extended: %{latitude: 35.5, longitude: -75.5}}
|
||||||
|
assert Index.extract_coordinate(packet, :lon) == -75.5
|
||||||
|
end
|
||||||
|
|
||||||
|
test "returns nil for missing data" do
|
||||||
|
assert Index.extract_coordinate(%{}, :lat) == nil
|
||||||
|
assert Index.extract_coordinate(%{}, :lon) == nil
|
||||||
|
end
|
||||||
|
|
||||||
|
test "prefers direct keys over data_extended" do
|
||||||
|
packet = %{lat: 10.0, lon: 20.0, data_extended: %{latitude: 30.0, longitude: 40.0}}
|
||||||
|
assert Index.extract_coordinate(packet, :lat) == 10.0
|
||||||
|
assert Index.extract_coordinate(packet, :lon) == 20.0
|
||||||
|
end
|
||||||
|
end
|
||||||
|
|
||||||
|
describe "format_coordinate/1" do
|
||||||
|
test "formats float to 6 decimal places" do
|
||||||
|
assert Index.format_coordinate(35.123456789) == "35.123457"
|
||||||
|
end
|
||||||
|
|
||||||
|
test "formats integer" do
|
||||||
|
assert Index.format_coordinate(35) == "35"
|
||||||
|
end
|
||||||
|
|
||||||
|
test "truncates long binary coordinate" do
|
||||||
|
assert Index.format_coordinate("35.12345678901") == "35.123456"
|
||||||
|
end
|
||||||
|
|
||||||
|
test "passes through short binary" do
|
||||||
|
assert Index.format_coordinate("35.12") == "35.12"
|
||||||
|
end
|
||||||
|
|
||||||
|
test "returns empty string for nil" do
|
||||||
|
assert Index.format_coordinate(nil) == ""
|
||||||
|
end
|
||||||
|
end
|
||||||
|
end
|
||||||
Loading…
Add table
Reference in a new issue