diff --git a/TODO.md b/TODO.md index 575fe5c..b09f04d 100644 --- a/TODO.md +++ b/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] 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 -- [ ] Cluster packet distribution race — `PacketDistributor.distribute_packet/1` only broadcasts if currently leader; leadership change between receipt and broadcast drops packets ## 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` -- [ ] 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 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 - [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] 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 - [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] 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. -- [ ] 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 diff --git a/assets/js/app.js b/assets/js/app.js index 55cc233..06c394c 100644 --- a/assets/js/app.js +++ b/assets/js/app.js @@ -1,5 +1,3 @@ -console.log("app.js loading..."); - // If you want to use Phoenix channels, run `mix help phx.gen.channel` // to get started and then uncomment the line below. // import "./user_socket.js" @@ -86,10 +84,8 @@ const originalMapMounted = MapAPRSMap.mounted; Hooks.APRSMap = { ...MapAPRSMap, mounted() { - console.log("APRSMap wrapper mounted() called"); const self = this; loadMapBundle(() => { - console.log("Map bundle ready, calling original mounted"); if (originalMapMounted) { 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, { longPollFallbackMs: 2500, 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 window.addEventListener("phx:reconnect", (e) => { const delay = e.detail.delay || 1000; - console.log(`[LiveSocket] Reconnecting in ${delay}ms due to connection draining...`); setTimeout(() => { // Disconnect and reconnect to potentially land on a different server liveSocket.disconnect(); @@ -216,7 +210,6 @@ window.addEventListener("phx:live_socket:connect", (info) => { if (socket && socket.fallbackTimer) { clearTimeout(socket.fallbackTimer); socket.fallbackTimer = null; - console.log("[LiveSocket] Cleared fallback timer after successful connection"); } }); @@ -226,7 +219,6 @@ setTimeout(() => { if (socket && socket.isConnected() && socket.fallbackTimer) { clearTimeout(socket.fallbackTimer); socket.fallbackTimer = null; - console.log("[LiveSocket] Cleared lingering fallback timer"); } }, 5000); diff --git a/lib/aprsme/cluster/leader_election.ex b/lib/aprsme/cluster/leader_election.ex index 0889069..6ff47e2 100644 --- a/lib/aprsme/cluster/leader_election.ex +++ b/lib/aprsme/cluster/leader_election.ex @@ -20,6 +20,15 @@ defmodule Aprsme.Cluster.LeaderElection do GenServer.call(__MODULE__, :leader?) 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 GenServer.call(__MODULE__, :current_leader) end @@ -60,6 +69,9 @@ defmodule Aprsme.Cluster.LeaderElection do # Schedule periodic checks 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}} end @@ -117,6 +129,7 @@ defmodule Aprsme.Cluster.LeaderElection do case :global.register_name(@election_key, self(), &resolve_conflict/3) do :yes -> Logger.info("Elected as APRS-IS connection leader on node #{node()}") + :persistent_term.put({__MODULE__, :is_leader}, true) notify_leadership_change(true) {: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_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)}") + :persistent_term.put({__MODULE__, :is_leader}, false) {:noreply, %{state | is_leader: false, leader_node: leader_node}} end end @@ -161,6 +175,7 @@ defmodule Aprsme.Cluster.LeaderElection do def terminate(reason, state) do if state.is_leader do Logger.info("Leader stepping down due to: #{inspect(reason)}") + :persistent_term.put({__MODULE__, :is_leader}, false) :global.unregister_name(@election_key) notify_leadership_change(false) end diff --git a/lib/aprsme/cluster/packet_distributor.ex b/lib/aprsme/cluster/packet_distributor.ex index 00e5a79..6cc5044 100644 --- a/lib/aprsme/cluster/packet_distributor.ex +++ b/lib/aprsme/cluster/packet_distributor.ex @@ -7,23 +7,19 @@ defmodule Aprsme.Cluster.PacketDistributor do alias Aprsme.Cluster.LeaderElection alias AprsmeWeb.MapLive.PacketStore - require Logger - @pubsub_topic "cluster:packets" def distribute_packet(packet) do # Only distribute if clustering is enabled and we're the leader 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 Phoenix.PubSub.broadcast( Aprsme.PubSub, @pubsub_topic, {:distributed_packet, packet} ) - - Logger.debug("Distributed packet #{packet.raw} to cluster") end end @@ -38,6 +34,6 @@ defmodule Aprsme.Cluster.PacketDistributor do # Update packet store for LiveView PacketStore.store_packet(packet) - Logger.debug("Received distributed packet on node #{node()}") + :ok end end diff --git a/lib/aprsme/packet_producer.ex b/lib/aprsme/packet_producer.ex index 6055c15..4a8e9c1 100644 --- a/lib/aprsme/packet_producer.ex +++ b/lib/aprsme/packet_producer.ex @@ -2,6 +2,9 @@ defmodule Aprsme.PacketProducer do @moduledoc """ GenStage producer that handles incoming APRS packets and sends them to consumers 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 @@ -17,14 +20,15 @@ defmodule Aprsme.PacketProducer do @impl true 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 @impl true - def handle_demand(incoming_demand, %{demand: demand, buffer: buffer} = state) do - {events, remaining_buffer, remaining_demand} = dispatch_events(buffer, demand + incoming_demand) - remaining_size = state.buffer_size - length(events) - {:noreply, events, %{state | demand: remaining_demand, buffer: remaining_buffer, buffer_size: remaining_size}} + def handle_demand(incoming_demand, %{demand: demand} = state) do + {events, new_buffer, remaining_demand} = dispatch_events(state.buffer, demand + incoming_demand) + + {:noreply, events, + %{state | demand: remaining_demand, buffer: new_buffer, buffer_size: state.buffer_size - length(events)}} end @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 - {: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 - defp dispatch_events(buffer, demand) when demand > 0 and buffer != [] do - {events, remaining} = Enum.split(buffer, demand) - {events, remaining, demand - length(events)} + defp dispatch_events(buffer, demand) when demand > 0 do + if :queue.is_empty(buffer) do + {[], 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 defp dispatch_events(buffer, demand) do diff --git a/lib/aprsme/packets.ex b/lib/aprsme/packets.ex index 70eaf7b..4e91b5a 100644 --- a/lib/aprsme/packets.ex +++ b/lib/aprsme/packets.ex @@ -225,13 +225,6 @@ defmodule Aprsme.Packets do # Ensure data_extended is properly sanitized before insertion 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 {:ok, packet} -> # Invalidate cache for this packet's callsign @@ -502,8 +495,6 @@ defmodule Aprsme.Packets do def get_recent_packets(opts \\ %{}) do 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 hours_back = Map.get(opts, :hours_back, 24) time_ago = DateTime.add(DateTime.utc_now(), -hours_back * 3600, :second) @@ -544,9 +535,7 @@ defmodule Aprsme.Packets do |> offset(^offset) |> QueryBuilder.with_coordinates() - result = Repo.all(query) - Logger.debug("Packets.get_recent_packets returning #{length(result)} packets") - result + Repo.all(query) end @doc """ diff --git a/lib/aprsme_web/live/packets_live/index.ex b/lib/aprsme_web/live/packets_live/index.ex index fabc4e3..e7748b2 100644 --- a/lib/aprsme_web/live/packets_live/index.ex +++ b/lib/aprsme_web/live/packets_live/index.ex @@ -32,4 +32,45 @@ defmodule AprsmeWeb.PacketsLive.Index do socket = assign(socket, :packets, packets) {:noreply, socket} 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 diff --git a/lib/aprsme_web/live/packets_live/index.html.heex b/lib/aprsme_web/live/packets_live/index.html.heex index ca00088..ea3be91 100644 --- a/lib/aprsme_web/live/packets_live/index.html.heex +++ b/lib/aprsme_web/live/packets_live/index.html.heex @@ -64,72 +64,12 @@