From 753fa50463b672867ca069ba54b51ed823baeb93 Mon Sep 17 00:00:00 2001 From: Graham McIntire Date: Thu, 19 Feb 2026 18:18:14 -0600 Subject: [PATCH] Fix ConnectionMonitor deadlock and LeaderElection silent leadership loss ConnectionMonitor.gather_cluster_stats was calling :rpc.call(Node.self()) which GenServer.call'd back into itself while blocked in handle_info, causing a 5-second deadlock timeout every 30 seconds. Fixed by computing local stats directly from state. LeaderElection never verified that a leader still held the :global registration. After :global conflict resolution (two isolated nodes merging), the loser was silently unregistered but stayed is_leader: true forever, causing both pods to connect to APRS-IS with the same callsign in a reconnect loop. Fixed by adding verify_leadership/1 to check_leadership which detects lost registrations and steps down. --- TODO.md | 298 +++++++++++++++++++ lib/aprsme/cluster/leader_election.ex | 20 ++ lib/aprsme/connection_monitor.ex | 24 +- test/aprsme/cluster/leader_election_test.exs | 71 ++++- test/aprsme/connection_monitor_test.exs | 10 +- 5 files changed, 410 insertions(+), 13 deletions(-) create mode 100644 TODO.md diff --git a/TODO.md b/TODO.md new file mode 100644 index 0000000..e8be6aa --- /dev/null +++ b/TODO.md @@ -0,0 +1,298 @@ +# Code Review Findings + +## Critical Bugs + +### 1. LiveView: Missing `{:noreply, socket}` in drain handler — crashes LiveView +**File:** `lib/aprsme_web/live/map_live/index.ex:885-896` + +The `true` branch returns a bare socket instead of `{:noreply, socket}`, which will crash the LiveView with a `MatchError` when the drain condition is hit. + +```elixir +# BUG: true branch returns bare socket +def handle_info({:drain_connections, to_drain}, socket) do + if :rand.uniform(100) <= to_drain * 10 do + socket # <-- needs {:noreply, socket |> put_flash(...) |> push_event(...)} + |> put_flash(:info, "Server load balancing in progress. Reconnecting...") + |> push_event("reconnect", %{delay: :rand.uniform(5000)}) + else + {:noreply, socket} + end +end +``` + +### ~~2. ConnectionMonitor: Deadlock via self-RPC every 30 seconds~~ FIXED +**File:** `lib/aprsme/connection_monitor.ex:126-141` + +~~`gather_cluster_stats/1` iterates `[Node.self() | Node.list()]` and calls `:rpc.call(node, __MODULE__, :get_stats, [])` for the local node. Since `get_stats/0` does a `GenServer.call(__MODULE__, :get_stats)`, and `gather_cluster_stats` is called from within `handle_info`, the GenServer is blocked waiting for the RPC which is waiting for the GenServer. Classic deadlock — times out every 30s.~~ + +**Fixed:** Now computes local stats directly from state instead of RPC to self. Remote nodes still use RPC. + +### 3. LiveView: `get_callsign_key` generates random keys for Ecto structs — breaks all dedup +**File:** `lib/aprsme_web/live/shared/packet_utils.ex:13-14` + +```elixir +def get_callsign_key(%{"id" => id}), do: to_string(id) +def get_callsign_key(_packet), do: [:positive] |> System.unique_integer() |> to_string() +``` + +Only matches string key `"id"`. Ecto structs with atom key `:id` fall through and get a random integer key every time. This means: +- Dedup never works (same packet gets different keys each call) +- `visible_packets` map accumulates duplicates +- Cleanup comparison (lines 1756-1766) never finds matches — all packets marked as expired + +### 4. Cluster: `PacketDistributor` never calls `SpatialPubSub.broadcast_packet` +**File:** `lib/aprsme/cluster/packet_distributor.ex:30-37` + +`handle_distributed_packet` calls `StreamingPacketsPubSub.broadcast_packet` but never calls `SpatialPubSub.broadcast_packet`. Any SpatialPubSub subscriber receives zero packets when clustering is enabled. + +### 5. LeaderElection: `pid_alive?/2` treats `{:badrpc, _}` as truthy +**File:** `lib/aprsme/cluster/leader_election.ex:257-263` + +`:rpc.call/4` returns `{:badrpc, reason}` on failure. Since tuples are truthy in Elixir, the function returns `true` when the RPC fails, preventing cleanup of stale leader registrations. + +### ~~5b. LeaderElection: Loser never detects lost registration after `:global` conflict resolution~~ FIXED +**File:** `lib/aprsme/cluster/leader_election.ex` + +When two pods start in isolation (before the Erlang cluster forms), both register as leader in their own `:global` namespace. When the cluster merges, `:global` calls `resolve_conflict/3` and silently unregisters the loser — but the loser's `check_leadership` handler only re-attempted election for non-leaders, so it stayed `is_leader: true` forever with an active APRS-IS connection. This caused both pods to connect to APRS-IS with the same callsign, creating a reconnect loop where they kicked each other off every 5 seconds. + +**Fixed:** Added `verify_leadership/1` which checks `:global.whereis_name` on every `check_leadership` tick. If the registration no longer points to `self()`, the node steps down and notifies `ConnectionManager` to stop the APRS-IS connection. + +### 6. SpatialPubSub: `ensure_float/1` crashes on integer strings +**File:** `lib/aprsme/spatial_pubsub.ex:239` + +`String.to_float("42")` raises `ArgumentError`. If bounds arrive as integer strings, this crashes the singleton SpatialPubSub GenServer. + +--- + +## Dual Packet State System (Architectural Bug) + +### 7. LiveView maintains two unsynchronized packet tracking systems +**Files:** `lib/aprsme_web/live/map_live/index.ex`, `packet_processor.ex`, `display_manager.ex` + +**System A:** `packet_state` wrapping `PacketManager`/`PacketStore` (ETS-backed). +**System B:** `all_packets`, `visible_packets`, `historical_packets` as plain maps in assigns. + +These are never synchronized: +- Real-time packets update System B but never touch System A +- Cleanup runs on System A but never touches System B +- Bounds filtering reads System A but display reads System B +- `PacketManager.extract_packet_ids` uses different ID generation than `PacketStore.generate_packet_id` — stored packets can never be found through the manager API + +--- + +## Performance Issues + +### 8. N+1 query: `has_weather_packets?` per marker +**File:** `lib/aprsme_web/live/map_live/data_builder.ex:334` + +Every non-weather marker on the map triggers an individual `Repo.exists?` query. With hundreds of visible markers, this is hundreds of DB round-trips per map load. Should batch-query weather callsigns upfront. + +### 9. Duplicate packet transformation — every packet processed twice +**File:** `lib/aprsme/is/is.ex:596-604` and `lib/aprsme/packet_consumer.ex:382-430` + +`Is.dispatch/1` runs `struct_to_map`, `extract_additional_data`, and `normalize_data_type`. Then `PacketConsumer.prepare_packet_for_insert/1` runs them all again. Every packet does double work. + +### 10. `batch_length/1` recomputes O(n) list length on every event arrival +**File:** `lib/aprsme/packet_consumer.ex:64-68, 130` + +`batch_length(batch)` calls `length/1` (O(n)) on every `handle_events` call instead of tracking the count as an integer in state. + +### 11. `check_batch_utilization` calls `length/1` four times in guards and body +**File:** `lib/aprsme/packet_consumer.ex:174-180` + +`length/1` in a guard traverses the entire list. Called again 3 more times in the body. For 1000-element batches, this is 4 full traversals just for a warning log. + +### 12. `StreamingPacketsPubSub` does full ETS scan on every packet +**File:** `lib/aprsme/streaming_packets_pubsub.ex:108-115` + +`:ets.select` with no match conditions returns every subscriber for every packet. Geographic filtering happens afterward. With N subscribers and M packets/sec, this is O(N*M). A plain map in GenServer state would be equivalent. + +### 13. `SpatialPubSub` receives every packet twice in non-cluster mode +**File:** `lib/aprsme/spatial_pubsub.ex:66` and `lib/aprsme/packet_consumer.ex:378` + +Receives packets from both `PacketConsumer.broadcast_single_packet` (direct cast) and `PostgresNotifier` (via PubSub). Duplicate processing and potential duplicate delivery. + +### 14. `remove_markers_batch` sends one WebSocket event per marker +**File:** `lib/aprsme_web/live/map_live/display_manager.ex:120-124` + +Each marker removal is a separate `push_event`. For 100 markers, that's 100 separate WebSocket events. Should batch into a single `push_event("remove_markers_batch", %{ids: marker_ids})`. + +### 15. `Encoding.clean_control_characters` allocates per grapheme for every string field +**File:** `lib/aprsme/encoding.ex:147-153` + +Decomposes into grapheme clusters, filters, joins — even when the string is clean ASCII (common case). A fast-path check (`String.printable?/1`) before decomposition would avoid allocation. + +### 16. `truncate_datetimes_to_second` called twice per packet +**File:** `lib/aprsme/packet_consumer.ex:278` and `lib/aprsme/packet_consumer.ex:426` + +Every datetime traversed and pattern-matched twice through the same function. + +### 17. Spatial queries don't use PostGIS GiST index +**File:** `lib/aprsme/packets/query_builder.ex:139-149` + +`within_bounds` filters on `p.lat`/`p.lon` B-tree columns instead of using `ST_Within` or `&&` operator on the `location` geometry column with a GiST index. Similarly `PreparedQueries` uses `ST_Y(location)` / `ST_X(location)` with BETWEEN, which also prevents spatial index use. + +### 18. `weather_only` query uses 10-way OR — can't use indexes +**File:** `lib/aprsme/packets/query_builder.ex:78-94` + +10-way OR across 10 columns forces sequential scan. The `has_weather` boolean column exists on the Packet schema but is not used by this query. + +### 19. `SpatialPubSub.get_intersecting_grid_cells` generates 341 cells for date-line-crossing viewports +**File:** `lib/aprsme/spatial_pubsub.ex:284-295` + +When west=170, east=-170, the range `170..-170` generates 341 cells instead of ~21. Functionally correct (filtered by `point_in_bounds?` later) but very wasteful. + +### 20. Cleanup O(n^2) packet comparison in LiveView +**File:** `lib/aprsme_web/live/map_live/index.ex:1756-1766` + +`Enum.reject` with `Enum.any?` inside = O(n*m). Combined with bug #3 (random keys), this always marks all packets as expired. + +--- + +## Memory Issues + +### 21. ETS cache has no TTL — unbounded growth +**File:** `lib/aprsme/cache.ex:20-27` + +`Cache.put/4` accepts TTL opts but ignores them. `ttl/2` always returns `{:ok, nil}`. `:query_cache`, `:device_cache`, and `:symbol_cache` grow without bound. + +### 22. Monitor reference leak in `StreamingPacketsPubSub` +**File:** `lib/aprsme/streaming_packets_pubsub.ex:77` + +Each `subscribe_to_bounds` creates a new monitor but the ref is never stored. `update_subscription_bounds` creates new monitors without demonitoring old ones. Accumulates dangling monitors per client session. + +### 23. `all_packets` map in LiveView assigns holds 2000 full Ecto structs +**File:** `lib/aprsme_web/live/map_live/packet_processor.ex:27-37` + +Each packet has 100+ columns. 2000 full structs per connected client. Never cleaned by `handle_cleanup_old_packets` (which only touches `packet_state`). + +### 24. `broadcast_async` uses `async_nolink` but never awaits — orphaned reply messages +**File:** `lib/aprsme/broadcast_task_supervisor.ex:47-58` + +`async_nolink` sends `{ref, result}` and `{:DOWN, ...}` messages to the calling GenServer. These pile up in the mailbox. Should use `start_child` (fire-and-forget) instead. + +--- + +## Dead Code + +### 25. `handle_info({:ssl, ...})` in Is module — connection uses `:gen_tcp`, not SSL +**File:** `lib/aprsme/is/is.ex:318-352` + +### 26. `PacketPipelineSetup` module — not in supervision tree, references unnamed consumers +**File:** `lib/aprsme/packet_pipeline_setup.ex` + +### 27. `Archiver` module — subscribes to PubSub but discards all messages +**File:** `lib/aprsme/archiver.ex` + +### 28. `SystemMonitor` — not in supervision tree, entirely inert +**File:** `lib/aprsme/system_monitor.ex` + +### 29. `DbOptimizer.copy_insert` — `COPY FROM STDIN` doesn't work through Postgrex +**File:** `lib/aprsme/db_optimizer.ex:23-39` + +Always fails and falls through to rescue which calls `regular_batch_insert`. Dead code path masked by rescue. + +### 30. Router `/health` route — shadowed by endpoint HealthCheck plug +**File:** `lib/aprsme_web/router.ex:56` vs `lib/aprsme_web/endpoint.ex:65` + +### 31. `DeviceIdentification.enqueue_refresh_job/0` — commented-out implementation +**File:** `lib/aprsme/device_identification.ex:192-194` + +### 32. `RateLimiterWrapper.count/2` and `reset/1` — stubs returning hardcoded values +**File:** `lib/aprsme/rate_limiter_wrapper.ex:16-29` + +### 33. `pubsub_config/0` branches return identical values +**File:** `lib/aprsme/application.ex:179-196` + +Both `if cluster_enabled` and `else` return `{Phoenix.PubSub, name: Aprsme.PubSub}`. + +### 34. Three no-op `send_heat_map_for_current_bounds` functions +**Files:** `lib/aprsme_web/live/map_live/packet_processor.ex:147`, `historical_loader.ex:414`, `display_manager.ex:127` + +Heat map never updates from real-time packets, historical loading at low zoom, or zoom threshold crossing. + +### 35. `placeholders` option passed to `Repo.insert_all` — not a valid option +**Files:** `lib/aprsme/packet_consumer.ex:304`, `lib/aprsme/db_optimizer.ex:64` + +--- + +## Correctness Issues + +### 36. Circuit breaker half-open state never persisted — allows unlimited concurrent probes +**File:** `lib/aprsme/circuit_breaker.ex:89-93` + +`calculate_current_state` returns `:half_open` but never writes it back to state. Multiple callers can all enter half-open simultaneously. + +### 37. `Process.cancel_timer(nil)` crash risk in Is module +**File:** `lib/aprsme/is/is.ex:321, 356` + +If `state.timer` is nil, `Process.cancel_timer(nil)` raises `FunctionClauseError`. Some paths protect with `if state.timer`, but the main data handlers don't. + +### 38. `PacketPipelineSupervisor` uses `:one_for_one` — producer restart orphans consumers +**File:** `lib/aprsme/packet_pipeline_supervisor.ex:23` + +If the producer crashes, consumers remain subscribed to the dead PID. Should use `:rest_for_one`. + +### 39. `handle_update_time_display` doesn't trigger re-renders +**File:** `lib/aprsme_web/live/map_live/index.ex:1774-1781` + +Returns `{:noreply, socket}` with unchanged socket — LiveView diff is empty. Time-ago displays never update. + +### 40. `PacketStore` uses global named ETS table shared across all LiveViews +**File:** `lib/aprsme_web/live/map_live/packet_store.ex:12, 114` + +Different LiveViews can overwrite each other's stored packets. TTL cleanup runs globally. + +### 41. ConnectionMonitor memory calculation is inverted +**File:** `lib/aprsme/connection_monitor.ex:219-232` + +`total / system` produces a ratio >1.0 (typically 1.5-3.0), not a percentage. + +### 42. `Callsign.valid?/1` accepts any non-empty string — docstring examples are wrong +**File:** `lib/aprsme/callsign.ex:40-45` + +Docstring says `valid?("A")` returns false, but implementation returns true. + +### 43. `Encoding.latin1_to_utf8` and `clean_control_characters` fight each other +**File:** `lib/aprsme/encoding.ex:126-170` + +`latin1_to_utf8` converts bytes 128-159 to Unicode U+0080-U+009F, then `clean_control_characters` immediately strips them. + +### 44. `DeviceCache` has unreachable outer `try/rescue` +**File:** `lib/aprsme/device_cache.ex:109-132` + +Inner rescue catches all exceptions (including the ones the outer rescue targets). Outer rescue for `Postgrex.Error`/`DBConnection.ConnectionError` is dead code. + +--- + +## Security / Operational + +### 45. LiveDashboard and ErrorTracker are publicly accessible +**File:** `lib/aprsme_web/router.ex:40-43` + +`/dashboard` and `/errors` are in a scope with only the `:browser` pipeline — no auth required. Exposes system metrics, process info, ETS tables, and stack traces. + +### 46. SQL injection in `DbOptimizer` table/column name interpolation +**File:** `lib/aprsme/db_optimizer.ex:29, 107, 126` + +`table_name` and `columns` interpolated directly into SQL. Internal-only callers today, but the pattern is dangerous. + +### 47. `/ready` and `/status.json` are rate-limited — risk of false health check failures +**File:** `lib/aprsme_web/router.ex:54-59` + +Piped through `:public_api` which includes `RateLimiter` at 100 req/min. Kubernetes probes hitting from same IP could get rate-limited. + +### 48. `ShutdownHandler.terminate` blocks 30s but scheduled messages never processed +**File:** `lib/aprsme/shutdown_handler.ex:69-92` + +`Process.send_after` in terminate is pointless — no more messages processed after terminate. The 30s sleep just delays shutdown. + +### 49. DatabaseMetrics hardcodes pool metrics +**File:** `lib/aprsme/telemetry/database_metrics.ex:77-85` + +Reports `busy: 2` and `idle: pool_size - 2` as actual telemetry. Misleads monitoring dashboards. + +### 50. `PacketBatcher` started outside supervision tree +**File:** `lib/aprsme_web/live/map_live/index.ex:202` + +If the batcher crashes, it won't restart. LiveView continues running but silently stops processing real-time packets. diff --git a/lib/aprsme/cluster/leader_election.ex b/lib/aprsme/cluster/leader_election.ex index 51a0c9b..728a71e 100644 --- a/lib/aprsme/cluster/leader_election.ex +++ b/lib/aprsme/cluster/leader_election.ex @@ -146,6 +146,8 @@ defmodule Aprsme.Cluster.LeaderElection do @impl true def handle_info(:check_leadership, state) do + state = verify_leadership(state) + # Re-attempt election if we're not leader if not state.is_leader do Process.send_after(self(), :attempt_election, 100) @@ -185,6 +187,24 @@ defmodule Aprsme.Cluster.LeaderElection do :ok end + # Verify that a node which thinks it's leader still holds the :global registration. + # After :global conflict resolution (e.g. two partitions merging), the losing PID + # is silently unregistered — this detects that and steps down. + defp verify_leadership(%{is_leader: true} = state) do + case :global.whereis_name(@election_key) do + pid when pid == self() -> + state + + _ -> + Logger.warning("Lost global leadership registration, stepping down") + :persistent_term.put({__MODULE__, :is_leader}, false) + notify_leadership_change(false) + %{state | is_leader: false, leader_node: nil} + end + end + + defp verify_leadership(state), do: state + defp attempt_registration(state) do case :global.register_name(@election_key, self(), &resolve_conflict/3) do :yes -> diff --git a/lib/aprsme/connection_monitor.ex b/lib/aprsme/connection_monitor.ex index 7ee49c4..a837d1e 100644 --- a/lib/aprsme/connection_monitor.ex +++ b/lib/aprsme/connection_monitor.ex @@ -124,20 +124,26 @@ defmodule Aprsme.ConnectionMonitor do end defp gather_cluster_stats(state) do - # Get stats from all nodes - nodes = [Node.self() | Node.list()] + # Compute local stats directly to avoid deadlock — calling get_stats via + # RPC on Node.self() would GenServer.call back into this process which is + # blocked in handle_info(:check_load, ...). + local_stats = %{ + connections: state.local_connections, + cpu: get_cpu_usage(), + memory: get_memory_usage(), + draining: state.draining + } - node_stats = - Enum.reduce(nodes, %{}, fn node, acc -> - stats = :rpc.call(node, __MODULE__, :get_stats, []) - - case stats do + # Get stats from remote nodes via RPC + remote_stats = + Enum.reduce(Node.list(), %{}, fn remote_node, acc -> + case :rpc.call(remote_node, __MODULE__, :get_stats, []) do {:badrpc, _} -> acc - stats -> Map.put(acc, node, stats) + stats -> Map.put(acc, remote_node, stats) end end) - %{state | node_stats: node_stats} + %{state | node_stats: Map.put(remote_stats, Node.self(), local_stats)} end defp analyze_load(state) do diff --git a/test/aprsme/cluster/leader_election_test.exs b/test/aprsme/cluster/leader_election_test.exs index bf2d36d..c242156 100644 --- a/test/aprsme/cluster/leader_election_test.exs +++ b/test/aprsme/cluster/leader_election_test.exs @@ -19,7 +19,11 @@ defmodule Aprsme.Cluster.LeaderElectionTest do :global.unregister_name(@election_key) if Process.whereis(LeaderElection) do - GenServer.stop(LeaderElection) + try do + GenServer.stop(LeaderElection) + catch + :exit, _ -> :ok + end end end) @@ -252,6 +256,71 @@ defmodule Aprsme.Cluster.LeaderElectionTest do end end + describe "leadership verification" do + setup do + Application.put_env(:aprsme, :cluster_enabled, false) + {:ok, pid} = LeaderElection.start_link([]) + Process.sleep(200) + assert LeaderElection.leader?() == true + {:ok, pid: pid} + end + + test "steps down when global registration is lost", %{pid: pid} do + # Subscribe to leadership change notifications + Phoenix.PubSub.subscribe(Aprsme.PubSub, "cluster:leadership") + + # Simulate losing the global registration (as happens during :global conflict resolution) + :global.unregister_name(@election_key) + + # Force a leadership check, then use :sys.get_state to synchronously wait + # for the message to be processed (before the 100ms re-election timer fires) + send(pid, :check_leadership) + state = :sys.get_state(pid) + + # Should have detected the loss and stepped down + assert state.is_leader == false + assert LeaderElection.leader_cached?() == false + + # Should have notified about leadership loss + assert_received {:leadership_change, _, false} + end + + test "steps down when another process holds the registration", %{pid: pid} do + Phoenix.PubSub.subscribe(Aprsme.PubSub, "cluster:leadership") + + # Simulate another node winning the conflict resolution + imposter = spawn(fn -> Process.sleep(:infinity) end) + :global.re_register_name(@election_key, imposter) + + # Force a leadership check + send(pid, :check_leadership) + state = :sys.get_state(pid) + + # Should have stepped down + assert state.is_leader == false + assert LeaderElection.leader_cached?() == false + assert_received {:leadership_change, _, false} + + # Clean up + Process.exit(imposter, :kill) + end + + test "re-elects after stepping down", %{pid: pid} do + # Lose the registration + :global.unregister_name(@election_key) + + # Force check — should step down + send(pid, :check_leadership) + state = :sys.get_state(pid) + assert state.is_leader == false + + # check_leadership also schedules :attempt_election for non-leaders. + # Wait for re-election. + Process.sleep(300) + assert LeaderElection.leader?() == true + end + end + describe "message handling" do setup do Application.put_env(:aprsme, :cluster_enabled, false) diff --git a/test/aprsme/connection_monitor_test.exs b/test/aprsme/connection_monitor_test.exs index c532543..58b51b4 100644 --- a/test/aprsme/connection_monitor_test.exs +++ b/test/aprsme/connection_monitor_test.exs @@ -129,17 +129,21 @@ defmodule Aprsme.ConnectionMonitorTest do assert stats.draining == true end - test "gather_cluster_stats handles local node stats", %{pid: pid} do + test "gather_cluster_stats includes local node stats without deadlock", %{pid: pid} do ConnectionMonitor.register_connection() ConnectionMonitor.register_connection() - # Trigger check_load which calls gather_cluster_stats internally + # Trigger check_load which calls gather_cluster_stats internally. + # Before the fix, this would deadlock for 5s because gather_cluster_stats + # did an RPC call to Node.self() which GenServer.call'd back into itself. send(pid, :check_load) Process.sleep(100) state = :sys.get_state(pid) - # node_stats should have been populated assert is_map(state.node_stats) + # Local node stats must always be present + assert Map.has_key?(state.node_stats, Node.self()) + assert state.node_stats[Node.self()].connections == 2 end end end