Turned on :error_handling, :underspecs, and :unmatched_returns in mix.exs dialyzer config. The 97 warnings this surfaced were fixed in place rather than suppressed: - unmatched_return (79): explicit discard with `_ = ...` for fire-and-forget side effects (Process.cancel_timer, :ets.new, send/2), and pattern-matched `:ok = ...` for control-plane Phoenix.PubSub subscribe/unsubscribe/broadcast calls so a future return-shape change fails loud. - contract_supertype (18): tightened @spec arg and return types on data_builder, historical_loader, url_params, packet_utils, encoding_utils, aprs_symbol, weather_controller, packet_replay to match each function's actual success typing. No behavioural change. mix compile clean, 1008 tests pass, dialyzer count is now 0.
248 lines
6.4 KiB
Elixir
248 lines
6.4 KiB
Elixir
defmodule Aprsme.ConnectionMonitor do
|
|
@moduledoc """
|
|
Monitors system load and LiveView connections to implement connection draining
|
|
when load is imbalanced across cluster nodes.
|
|
"""
|
|
use GenServer
|
|
|
|
require Logger
|
|
|
|
@check_interval to_timeout(second: 30)
|
|
# Start draining if CPU usage > 70%
|
|
@cpu_threshold 0.7
|
|
# Drain if this node has 2x more connections than average
|
|
@connection_imbalance_ratio 2.0
|
|
# Drain 10% of connections when triggered
|
|
@drain_percentage 0.1
|
|
|
|
defstruct [
|
|
:node_stats,
|
|
:local_connections,
|
|
:draining,
|
|
:last_check
|
|
]
|
|
|
|
def start_link(opts) do
|
|
GenServer.start_link(__MODULE__, opts, name: __MODULE__)
|
|
end
|
|
|
|
@impl true
|
|
def init(_opts) do
|
|
# Only start if clustering is enabled
|
|
if cluster_enabled?() do
|
|
schedule_check()
|
|
|
|
{:ok,
|
|
%__MODULE__{
|
|
node_stats: %{},
|
|
local_connections: 0,
|
|
draining: false,
|
|
last_check: System.monotonic_time(:second)
|
|
}}
|
|
else
|
|
:ignore
|
|
end
|
|
end
|
|
|
|
@doc """
|
|
Register a new LiveView connection
|
|
"""
|
|
def register_connection do
|
|
if cluster_enabled?() do
|
|
GenServer.cast(__MODULE__, :register_connection)
|
|
end
|
|
end
|
|
|
|
@doc """
|
|
Unregister a LiveView connection
|
|
"""
|
|
def unregister_connection do
|
|
if cluster_enabled?() do
|
|
GenServer.cast(__MODULE__, :unregister_connection)
|
|
end
|
|
end
|
|
|
|
@doc """
|
|
Check if this node should accept new connections
|
|
"""
|
|
def accepting_connections? do
|
|
if cluster_enabled?() and Process.whereis(__MODULE__) != nil do
|
|
GenServer.call(__MODULE__, :accepting_connections?)
|
|
else
|
|
true
|
|
end
|
|
end
|
|
|
|
@doc """
|
|
Get current node statistics
|
|
"""
|
|
def get_stats do
|
|
if cluster_enabled?() and Process.whereis(__MODULE__) != nil do
|
|
GenServer.call(__MODULE__, :get_stats)
|
|
else
|
|
%{connections: 0, cpu: 0.0, memory: 0.0}
|
|
end
|
|
end
|
|
|
|
@impl true
|
|
def handle_cast(:register_connection, state) do
|
|
{:noreply, %{state | local_connections: state.local_connections + 1}}
|
|
end
|
|
|
|
@impl true
|
|
def handle_cast(:unregister_connection, state) do
|
|
{:noreply, %{state | local_connections: max(0, state.local_connections - 1)}}
|
|
end
|
|
|
|
@impl true
|
|
def handle_call(:accepting_connections?, _from, state) do
|
|
{:reply, not state.draining, state}
|
|
end
|
|
|
|
@impl true
|
|
def handle_call(:get_stats, _from, state) do
|
|
stats = %{
|
|
connections: state.local_connections,
|
|
cpu: get_cpu_usage(),
|
|
memory: get_memory_usage(),
|
|
draining: state.draining
|
|
}
|
|
|
|
{:reply, stats, state}
|
|
end
|
|
|
|
@impl true
|
|
def handle_info(:check_load, state) do
|
|
new_state =
|
|
state
|
|
|> gather_cluster_stats()
|
|
|> analyze_load()
|
|
|> maybe_trigger_draining()
|
|
|
|
schedule_check()
|
|
{:noreply, new_state}
|
|
end
|
|
|
|
defp gather_cluster_stats(state) do
|
|
# 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
|
|
}
|
|
|
|
# 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, remote_node, stats)
|
|
end
|
|
end)
|
|
|
|
%{state | node_stats: Map.put(remote_stats, Node.self(), local_stats)}
|
|
end
|
|
|
|
defp analyze_load(state) do
|
|
local_stats = Map.get(state.node_stats, Node.self(), %{})
|
|
cpu_usage = Map.get(local_stats, :cpu, 0.0)
|
|
|
|
# Calculate average connections across cluster
|
|
total_connections = Enum.sum(for {_, stats} <- state.node_stats, do: Map.get(stats, :connections, 0))
|
|
node_count = map_size(state.node_stats)
|
|
avg_connections = if node_count > 0, do: total_connections / node_count, else: 0
|
|
|
|
# Determine if we should be draining
|
|
should_drain =
|
|
cond do
|
|
# High CPU usage
|
|
cpu_usage > @cpu_threshold ->
|
|
Logger.info("Node #{Node.self()} CPU usage high: #{Float.round(cpu_usage * 100, 1)}%")
|
|
true
|
|
|
|
# Too many connections compared to average
|
|
node_count > 1 and state.local_connections > avg_connections * @connection_imbalance_ratio ->
|
|
Logger.info(
|
|
"Node #{Node.self()} has #{state.local_connections} connections, avg: #{Float.round(avg_connections, 1)}"
|
|
)
|
|
|
|
true
|
|
|
|
# Otherwise, stop draining if we were
|
|
true ->
|
|
false
|
|
end
|
|
|
|
%{state | draining: should_drain}
|
|
end
|
|
|
|
defp maybe_trigger_draining(%{draining: true, local_connections: connections} = state) when connections > 0 do
|
|
# Calculate how many connections to drain
|
|
to_drain = round(connections * @drain_percentage)
|
|
# Drain 1-10 connections at a time
|
|
to_drain = max(1, min(to_drain, 10))
|
|
|
|
Logger.info("Draining #{to_drain} connections from node #{Node.self()}")
|
|
|
|
# Broadcast drain event to LiveViews
|
|
_ =
|
|
Phoenix.PubSub.broadcast(
|
|
Aprsme.PubSub,
|
|
"connection:drain:#{Node.self()}",
|
|
{:drain_connections, to_drain}
|
|
)
|
|
|
|
state
|
|
end
|
|
|
|
defp maybe_trigger_draining(state), do: state
|
|
|
|
defp get_cpu_usage do
|
|
# Get CPU usage from scheduler utilization
|
|
# Use Erlang's cpu_sup if available, otherwise estimate from scheduler utilization
|
|
case :cpu_sup.util() do
|
|
util when is_number(util) ->
|
|
util / 100.0
|
|
|
|
_ ->
|
|
# Fallback to scheduler wall time
|
|
:scheduler_wall_time_all
|
|
|> :erlang.statistics()
|
|
|> Enum.map(&calculate_scheduler_utilization/1)
|
|
|> Enum.sum()
|
|
|> Kernel./(System.schedulers_online())
|
|
end
|
|
rescue
|
|
_ -> 0.0
|
|
end
|
|
|
|
defp calculate_scheduler_utilization({_, active, total}) do
|
|
if total > 0, do: active / total, else: 0.0
|
|
end
|
|
|
|
defp get_memory_usage do
|
|
mem_data = :erlang.memory()
|
|
total = Keyword.get(mem_data, :total, 0)
|
|
# processes + system = total; report process memory as fraction of total
|
|
processes = Keyword.get(mem_data, :processes, 0)
|
|
|
|
if total > 0 do
|
|
processes / total
|
|
else
|
|
0.0
|
|
end
|
|
rescue
|
|
_ -> 0.0
|
|
end
|
|
|
|
defp schedule_check do
|
|
Process.send_after(self(), :check_load, @check_interval)
|
|
end
|
|
|
|
defp cluster_enabled? do
|
|
Application.get_env(:aprsme, :cluster_enabled, false)
|
|
end
|
|
end
|