Packages

macula

0.20.21
7.1.0 7.0.0 6.0.0 5.2.2 5.2.1 5.2.0 5.1.0 5.0.0 4.8.0 4.7.1 4.7.0 4.6.0 4.5.0 4.4.10 4.4.9 4.4.8 4.4.7 4.4.6 4.4.5 4.4.4 4.4.3 4.4.2 4.4.1 4.4.0 4.3.1 4.3.0 4.2.9 4.2.8 4.2.7 4.2.6 4.2.5 4.2.4 4.2.3 4.2.2 4.2.1 4.2.0 4.1.1 4.1.0 4.0.0 3.16.0 3.15.3 3.15.2 3.15.1 3.14.0 3.13.0 3.12.1 3.12.0 3.11.1 3.11.0 3.10.3 3.10.2 3.10.1 3.9.0 3.8.0 3.7.0 3.5.0 3.4.0 3.3.0 3.2.0 3.1.0 3.0.0 2.1.1 2.1.0 2.0.0 1.5.2 1.5.1 1.4.30 1.4.29 1.4.28 1.4.27 1.4.26 1.4.25 1.4.24 1.4.23 1.4.22 1.4.21 1.4.20 1.4.19 1.4.18 1.4.17 1.4.16 1.4.15 1.4.14 1.4.13 1.4.11 1.4.10 1.4.9 1.4.8 1.4.7 1.4.6 1.4.5 1.4.4 1.4.3 1.4.2 1.4.1 1.4.0 1.3.1 1.3.0 1.2.0 1.1.0 1.0.10 1.0.9 1.0.8 1.0.7 1.0.6 1.0.5 1.0.4 1.0.3 1.0.2 1.0.1 1.0.0 0.48.6 0.48.5 0.48.4 0.48.3 0.48.2 0.48.1 0.48.0 0.47.1 0.47.0 0.46.3 0.46.1 0.46.0 0.45.3 0.45.2 0.45.1 0.45.0 0.44.2 0.44.1 0.44.0 0.43.3 0.43.2 0.43.1 0.43.0 0.42.9 0.42.8 0.42.7 0.42.6 0.42.5 0.42.4 0.42.3 0.42.2 0.42.1 0.42.0 0.41.1 0.41.0 0.40.1 0.40.0 0.39.9 0.39.8 0.39.7 0.39.6 0.39.5 0.39.4 0.39.3 0.39.2 0.39.1 0.39.0 0.38.8 0.38.7 0.38.6 0.38.5 0.38.4 0.38.3 0.38.2 0.38.1 0.38.0 0.37.7 0.37.6 0.37.5 0.37.4 0.37.3 0.37.2 0.37.1 0.37.0 0.36.6 0.36.5 0.36.4 0.36.3 0.36.2 0.36.1 0.36.0 0.35.4 0.35.3 0.35.2 0.35.1 0.35.0 0.34.1 0.34.0 0.33.1 0.33.0 0.32.5 0.32.4 0.32.3 0.32.2 0.32.1 0.32.0 0.31.9 0.31.8 0.31.7 0.31.6 0.31.5 0.31.4 0.31.3 0.31.2 0.31.1 0.31.0 0.30.10 0.30.9 0.30.8 0.30.7 0.30.6 0.30.5 0.30.4 0.30.3 0.30.2 0.30.1 0.30.0 0.29.0 0.28.3 0.28.2 0.28.1 0.28.0 0.27.1 0.27.0 0.26.1 0.26.0 0.25.6 0.25.5 0.25.4 0.25.3 0.25.2 0.25.1 0.25.0 0.24.6 0.24.5 0.24.4 0.24.3 0.24.2 0.24.1 0.24.0 0.23.3 0.23.2 0.23.1 0.23.0 0.22.12 0.22.11 0.22.10 0.22.9 0.22.8 0.22.7 0.22.6 0.22.5 0.22.4 0.22.3 0.22.2 0.22.1 0.22.0 0.21.7 0.21.6 0.21.5 0.21.4 0.21.2 0.21.1 0.21.0 0.20.25 0.20.24 0.20.23 0.20.22 0.20.21 0.20.20 0.20.19 0.20.18 0.20.17 0.20.16 0.20.15 0.20.14 0.20.13 0.20.12 0.20.11 0.20.10 0.20.9 0.20.8 0.20.7 0.20.6 0.20.5 0.20.3 0.20.2 0.20.1 0.20.0 0.19.2 0.19.1 0.19.0 0.18.1 0.18.0 0.17.4 0.17.3 0.17.2 0.17.1 0.17.0 0.16.6 0.16.5 0.16.4 0.16.3 0.16.2 0.16.1 0.16.0 0.15.1 0.15.0 0.14.3 0.14.2 0.14.1 0.14.0 0.12.6 0.12.5 0.12.3 0.11.3 0.10.2 0.10.1 0.10.0 0.9.2 0.9.1 0.9.0 0.8.25 0.8.24 0.8.23 0.8.22 0.8.21 0.8.20 0.8.19 0.8.18 0.8.17 0.8.16 0.8.15 0.8.14 0.8.13 0.8.12 0.8.11 0.8.10 0.8.9 0.8.8 0.8.7 0.8.6 0.8.5 0.8.4 0.8.3 0.8.2 0.8.1 0.8.0 0.7.30 0.7.29 0.7.28 0.7.27 0.7.26 0.7.25 0.7.24 0.7.23 0.7.22 0.7.21 0.7.20 0.7.19 0.7.18 0.7.17 0.7.16 0.7.15 0.7.14 0.7.13 0.7.12 0.7.11 0.7.10 0.7.9 0.7.8 0.7.7 0.7.6 0.7.5 0.7.4 0.7.3 0.7.2 0.7.1 0.7.0 0.6.7 0.6.6 0.6.5 0.6.4 0.6.3 0.6.2 0.6.1 0.6.0 0.5.0 0.4.4 0.4.3 0.4.2 0.4.1 0.4.0 0.3.4 0.3.3 0.3.2 0.3.1

Macula HTTP/3 Mesh SDK — connect, subscribe, publish, call, advertise

Current section

Files

Jump to
macula src macula_nat_system macula_port_predictor.erl
Raw

src/macula_nat_system/macula_port_predictor.erl

%%%-------------------------------------------------------------------
%%% @doc
%%% Port Prediction for NAT Traversal.
%%%
%%% Provides intelligent port prediction based on NAT allocation policies
%%% and historical port allocation data. Improves hole punch success rates
%%% by predicting the external ports a peer will use.
%%%
%%% Prediction Strategies by Allocation Policy:
%%%
%%% - PP (Port Preservation): External port = internal port
%%% Prediction: Use last known port (high confidence)
%%%
%%% - PC (Port Contiguity): Sequential port allocation
%%% Prediction: Calculate delta from history, predict next ports
%%% Delta tracking with exponential moving average
%%%
%%% - RD (Random): Random port allocation
%%% Prediction: Use historical range + common NAT port ranges
%%% Statistical approach with lower confidence
%%%
%%% Port History Storage:
%%% - Tracks last N port allocations per peer (default: 10)
%%% - Calculates port deltas for PC allocation
%%% - Maintains statistics for RD allocation (mean, stddev, range)
%%%
%%% @end
%%%-------------------------------------------------------------------
-module(macula_port_predictor).
-behaviour(gen_server).
-include_lib("kernel/include/logger.hrl").
%% API
-export([
start_link/1,
predict/2,
predict/3,
record_port/3,
get_history/1,
get_stats/1,
clear/1
]).
%% gen_server callbacks
-export([init/1, handle_call/3, handle_cast/2, handle_info/2, terminate/2]).
-define(SERVER, ?MODULE).
-define(TABLE, macula_port_history).
-define(DEFAULT_HISTORY_SIZE, 10).
-define(DEFAULT_PREDICTION_COUNT, 5).
-define(MAX_PORT, 65535).
-define(MIN_PORT, 1024).
-define(COMMON_NAT_PORT_RANGE_START, 32768).
-define(COMMON_NAT_PORT_RANGE_END, 61000).
-define(CLEANUP_INTERVAL_MS, 300000). % Cleanup every 5 minutes
-define(HISTORY_TTL_SECONDS, 1800). % 30 minutes TTL for history
%%%===================================================================
%%% Types
%%%===================================================================
-type allocation_policy() :: pp | pc | rd | unknown.
-type port_prediction() :: #{
ports := [inet:port_number()],
confidence := float(), % 0.0 - 1.0
strategy := pp | pc | rd | fallback,
delta => integer(), % For PC strategy
stats => port_stats() % For RD strategy
}.
-type port_history() :: #{
node_id := binary(),
ports := [inet:port_number()],
deltas := [integer()],
updated_at := integer()
}.
-type port_stats() :: #{
mean := float(),
stddev := float(),
min := inet:port_number(),
max := inet:port_number(),
count := non_neg_integer()
}.
-export_type([port_prediction/0, port_history/0, port_stats/0]).
-record(state, {
table :: ets:tid(),
history_size :: pos_integer(),
prediction_count :: pos_integer()
}).
%%%===================================================================
%%% API
%%%===================================================================
%% @doc Start the port predictor server.
-spec start_link(map()) -> {ok, pid()} | {error, term()}.
start_link(Opts) ->
gen_server:start_link({local, ?SERVER}, ?MODULE, Opts, []).
%% @doc Predict external ports for a peer.
%% Uses NAT profile if available, otherwise uses historical data.
-spec predict(binary(), allocation_policy()) -> port_prediction().
predict(NodeId, AllocationPolicy) ->
predict(NodeId, AllocationPolicy, #{}).
%% @doc Predict external ports with options.
%% Options:
%% base_port - Known/expected base port for prediction
%% count - Number of ports to predict (default: 5)
-spec predict(binary(), allocation_policy(), map()) -> port_prediction().
predict(NodeId, AllocationPolicy, Opts) ->
gen_server:call(?SERVER, {predict, NodeId, AllocationPolicy, Opts}).
%% @doc Record an observed external port for a peer.
%% Used to build history for better predictions.
-spec record_port(binary(), inet:port_number(), allocation_policy()) -> ok.
record_port(NodeId, Port, AllocationPolicy) ->
gen_server:cast(?SERVER, {record_port, NodeId, Port, AllocationPolicy}).
%% @doc Get port history for a peer.
-spec get_history(binary()) -> {ok, port_history()} | not_found.
get_history(NodeId) ->
gen_server:call(?SERVER, {get_history, NodeId}).
%% @doc Get port statistics for a peer.
-spec get_stats(binary()) -> {ok, port_stats()} | not_found.
get_stats(NodeId) ->
gen_server:call(?SERVER, {get_stats, NodeId}).
%% @doc Clear port history for a peer.
-spec clear(binary()) -> ok.
clear(NodeId) ->
gen_server:cast(?SERVER, {clear, NodeId}).
%%%===================================================================
%%% gen_server callbacks
%%%===================================================================
init(Opts) ->
HistorySize = maps:get(port_history_size, Opts, ?DEFAULT_HISTORY_SIZE),
PredictionCount = maps:get(port_prediction_count, Opts, ?DEFAULT_PREDICTION_COUNT),
%% Create ETS table for port history
Table = ets:new(?TABLE, [
set,
protected,
{keypos, 1},
{read_concurrency, true}
]),
%% Schedule periodic cleanup
schedule_cleanup(),
?LOG_INFO("Port predictor started (history_size=~p, prediction_count=~p)",
[HistorySize, PredictionCount]),
{ok, #state{
table = Table,
history_size = HistorySize,
prediction_count = PredictionCount
}}.
handle_call({predict, NodeId, AllocationPolicy, Opts}, _From, State) ->
#state{table = Table, prediction_count = DefaultCount} = State,
Count = maps:get(count, Opts, DefaultCount),
BasePort = maps:get(base_port, Opts, undefined),
Prediction = do_predict(Table, NodeId, AllocationPolicy, BasePort, Count),
{reply, Prediction, State};
handle_call({get_history, NodeId}, _From, State) ->
Result = case ets:lookup(State#state.table, NodeId) of
[{NodeId, History}] -> {ok, History};
[] -> not_found
end,
{reply, Result, State};
handle_call({get_stats, NodeId}, _From, State) ->
Result = case ets:lookup(State#state.table, NodeId) of
[{NodeId, #{ports := Ports}}] when length(Ports) > 0 ->
{ok, calculate_stats(Ports)};
_ ->
not_found
end,
{reply, Result, State};
handle_call(_Request, _From, State) ->
{reply, {error, unknown_request}, State}.
handle_cast({record_port, NodeId, Port, AllocationPolicy}, State) ->
#state{table = Table, history_size = MaxHistory} = State,
Now = erlang:system_time(second),
History = case ets:lookup(Table, NodeId) of
[{NodeId, Existing}] ->
update_history(Existing, Port, AllocationPolicy, MaxHistory, Now);
[] ->
new_history(NodeId, Port, AllocationPolicy, Now)
end,
ets:insert(Table, {NodeId, History}),
?LOG_DEBUG("Recorded port ~p for ~s (policy: ~p)", [Port, NodeId, AllocationPolicy]),
{noreply, State};
handle_cast({clear, NodeId}, State) ->
ets:delete(State#state.table, NodeId),
{noreply, State};
handle_cast(_Msg, State) ->
{noreply, State}.
handle_info(cleanup, State) ->
cleanup_expired_history(State#state.table),
schedule_cleanup(),
{noreply, State};
handle_info(_Info, State) ->
{noreply, State}.
terminate(_Reason, #state{table = Table}) ->
ets:delete(Table),
ok.
%%%===================================================================
%%% Internal functions - Prediction
%%%===================================================================
%% @private
%% @doc Perform port prediction based on allocation policy and history.
-spec do_predict(ets:tid(), binary(), allocation_policy(),
inet:port_number() | undefined, pos_integer()) -> port_prediction().
do_predict(Table, NodeId, pp, BasePort, _Count) ->
%% Port Preservation - use base port or last known port
predict_pp(Table, NodeId, BasePort);
do_predict(Table, NodeId, pc, BasePort, Count) ->
%% Port Contiguity - calculate delta and predict sequence
predict_pc(Table, NodeId, BasePort, Count);
do_predict(Table, NodeId, rd, BasePort, Count) ->
%% Random - use statistical approach
predict_rd(Table, NodeId, BasePort, Count);
do_predict(Table, NodeId, unknown, BasePort, Count) ->
%% Unknown - try PC first, fall back to range
predict_unknown(Table, NodeId, BasePort, Count).
%% @private
%% @doc Predict ports for PP (Port Preservation) NAT.
-spec predict_pp(ets:tid(), binary(), inet:port_number() | undefined) -> port_prediction().
predict_pp(_Table, _NodeId, BasePort) when is_integer(BasePort), BasePort > 0 ->
#{
ports => [BasePort],
confidence => 0.95,
strategy => pp
};
predict_pp(Table, NodeId, undefined) ->
case ets:lookup(Table, NodeId) of
[{NodeId, #{ports := [LastPort | _]}}] ->
#{
ports => [LastPort],
confidence => 0.85,
strategy => pp
};
_ ->
%% No history - can't predict PP
fallback_prediction()
end.
%% @private
%% @doc Predict ports for PC (Port Contiguity) NAT.
-spec predict_pc(ets:tid(), binary(), inet:port_number() | undefined, pos_integer()) ->
port_prediction().
predict_pc(Table, NodeId, BasePort, Count) ->
case get_average_delta(Table, NodeId) of
{ok, Delta} when Delta =/= 0 ->
Base = determine_base_port_pc(Table, NodeId, BasePort),
Ports = generate_sequential_ports(Base, Delta, Count),
#{
ports => Ports,
confidence => calculate_pc_confidence(Table, NodeId),
strategy => pc,
delta => Delta
};
_ ->
%% No reliable delta - use fallback with base port
case BasePort of
undefined -> fallback_prediction();
_ -> predict_around_port(BasePort, Count)
end
end.
%% @private
%% @doc Predict ports for RD (Random) NAT.
-spec predict_rd(ets:tid(), binary(), inet:port_number() | undefined, pos_integer()) ->
port_prediction().
predict_rd(Table, NodeId, BasePort, Count) ->
case ets:lookup(Table, NodeId) of
[{NodeId, #{ports := Ports}}] when length(Ports) >= 3 ->
Stats = calculate_stats(Ports),
PredictedPorts = predict_from_stats(Stats, BasePort, Count),
#{
ports => PredictedPorts,
confidence => calculate_rd_confidence(Stats),
strategy => rd,
stats => Stats
};
_ ->
%% Not enough history - use common NAT range
predict_common_nat_range(BasePort, Count)
end.
%% @private
%% @doc Predict ports when allocation policy is unknown.
-spec predict_unknown(ets:tid(), binary(), inet:port_number() | undefined, pos_integer()) ->
port_prediction().
predict_unknown(Table, NodeId, BasePort, Count) ->
%% Try to infer policy from history
case infer_allocation_policy(Table, NodeId) of
{ok, pp} -> predict_pp(Table, NodeId, BasePort);
{ok, pc} -> predict_pc(Table, NodeId, BasePort, Count);
{ok, rd} -> predict_rd(Table, NodeId, BasePort, Count);
unknown ->
case BasePort of
undefined -> fallback_prediction();
_ -> predict_around_port(BasePort, Count)
end
end.
%% @private
%% @doc Generate fallback prediction when no data available.
-spec fallback_prediction() -> port_prediction().
fallback_prediction() ->
%% Use common NAT ephemeral port range
Ports = [
?COMMON_NAT_PORT_RANGE_START,
?COMMON_NAT_PORT_RANGE_START + 1000,
?COMMON_NAT_PORT_RANGE_START + 2000,
?COMMON_NAT_PORT_RANGE_START + 3000,
?COMMON_NAT_PORT_RANGE_START + 4000
],
#{
ports => Ports,
confidence => 0.1,
strategy => fallback
}.
%% @private
%% @doc Predict ports around a known base port.
-spec predict_around_port(inet:port_number(), pos_integer()) -> port_prediction().
predict_around_port(BasePort, Count) ->
HalfCount = Count div 2,
Offsets = lists:seq(-HalfCount, HalfCount),
Ports = [clamp_port(BasePort + Offset) || Offset <- Offsets],
UniquePorts = lists:usort(Ports),
#{
ports => lists:sublist(UniquePorts, Count),
confidence => 0.5,
strategy => fallback
}.
%% @private
%% @doc Predict ports in common NAT range.
-spec predict_common_nat_range(inet:port_number() | undefined, pos_integer()) ->
port_prediction().
predict_common_nat_range(BasePort, Count) when is_integer(BasePort) ->
%% Center predictions around base port within NAT range
RangeSize = ?COMMON_NAT_PORT_RANGE_END - ?COMMON_NAT_PORT_RANGE_START,
Step = RangeSize div (Count + 1),
BasePorts = [clamp_port(BasePort + (I - Count div 2) * Step) || I <- lists:seq(0, Count - 1)],
#{
ports => BasePorts,
confidence => 0.2,
strategy => rd
};
predict_common_nat_range(undefined, Count) ->
%% Spread across common NAT range
RangeSize = ?COMMON_NAT_PORT_RANGE_END - ?COMMON_NAT_PORT_RANGE_START,
Step = RangeSize div (Count + 1),
Ports = [?COMMON_NAT_PORT_RANGE_START + (I * Step) || I <- lists:seq(1, Count)],
#{
ports => Ports,
confidence => 0.15,
strategy => rd
}.
%%%===================================================================
%%% Internal functions - History Management
%%%===================================================================
%% @private
%% @doc Create new port history entry.
-spec new_history(binary(), inet:port_number(), allocation_policy(), integer()) ->
port_history().
new_history(NodeId, Port, _AllocationPolicy, Now) ->
#{
node_id => NodeId,
ports => [Port],
deltas => [],
updated_at => Now
}.
%% @private
%% @doc Update existing port history.
-spec update_history(port_history(), inet:port_number(), allocation_policy(),
pos_integer(), integer()) -> port_history().
update_history(#{ports := Ports, deltas := Deltas} = History, Port, _Policy, MaxHistory, Now) ->
NewPorts = lists:sublist([Port | Ports], MaxHistory),
%% Calculate delta from previous port
NewDeltas = case Ports of
[PrevPort | _] ->
Delta = Port - PrevPort,
lists:sublist([Delta | Deltas], MaxHistory - 1);
[] ->
Deltas
end,
History#{
ports := NewPorts,
deltas := NewDeltas,
updated_at := Now
}.
%% @private
%% @doc Get average delta from port history.
-spec get_average_delta(ets:tid(), binary()) -> {ok, integer()} | not_found.
get_average_delta(Table, NodeId) ->
case ets:lookup(Table, NodeId) of
[{NodeId, #{deltas := Deltas}}] when length(Deltas) >= 2 ->
%% Use exponential weighted moving average (recent deltas weighted more)
AvgDelta = calculate_weighted_average(Deltas),
{ok, round(AvgDelta)};
_ ->
not_found
end.
%% @private
%% @doc Calculate exponentially weighted moving average.
-spec calculate_weighted_average([integer()]) -> float().
calculate_weighted_average([]) -> 0.0;
calculate_weighted_average(Deltas) ->
Alpha = 0.3, % Weight for most recent value
calculate_ewma(Deltas, Alpha, 0.0, 0.0).
calculate_ewma([], _Alpha, Sum, TotalWeight) when TotalWeight > 0 ->
Sum / TotalWeight;
calculate_ewma([], _Alpha, _Sum, _TotalWeight) ->
0.0;
calculate_ewma([Delta | Rest], Alpha, Sum, TotalWeight) ->
Weight = math:pow(1 - Alpha, length(Rest)),
calculate_ewma(Rest, Alpha, Sum + (Delta * Weight), TotalWeight + Weight).
%% @private
%% @doc Determine base port for PC prediction.
-spec determine_base_port_pc(ets:tid(), binary(), inet:port_number() | undefined) ->
inet:port_number().
determine_base_port_pc(_Table, _NodeId, BasePort) when is_integer(BasePort), BasePort > 0 ->
BasePort;
determine_base_port_pc(Table, NodeId, undefined) ->
case ets:lookup(Table, NodeId) of
[{NodeId, #{ports := [LastPort | _]}}] -> LastPort;
_ -> ?COMMON_NAT_PORT_RANGE_START
end.
%% @private
%% @doc Generate sequential ports based on delta.
-spec generate_sequential_ports(inet:port_number(), integer(), pos_integer()) ->
[inet:port_number()].
generate_sequential_ports(Base, Delta, Count) ->
HalfCount = Count div 2,
%% Generate ports both forward and backward from base
Offsets = lists:seq(-HalfCount, HalfCount + (Count rem 2)),
Ports = [clamp_port(Base + (Offset * Delta)) || Offset <- Offsets],
lists:usort(Ports).
%% @private
%% @doc Calculate confidence for PC prediction.
-spec calculate_pc_confidence(ets:tid(), binary()) -> float().
calculate_pc_confidence(Table, NodeId) ->
case ets:lookup(Table, NodeId) of
[{NodeId, #{deltas := Deltas}}] when length(Deltas) >= 3 ->
%% Confidence based on delta consistency (low variance = high confidence)
Variance = calculate_variance(Deltas),
%% Map variance to confidence: low variance = high confidence
max(0.3, min(0.9, 1.0 - (Variance / 1000)));
[{NodeId, #{deltas := Deltas}}] when length(Deltas) >= 1 ->
%% Some history - moderate confidence
0.5;
_ ->
0.3
end.
%%%===================================================================
%%% Internal functions - Statistics
%%%===================================================================
%% @private
%% @doc Calculate statistics from port list.
-spec calculate_stats([inet:port_number()]) -> port_stats().
calculate_stats([]) ->
#{mean => 0.0, stddev => 0.0, min => 0, max => 0, count => 0};
calculate_stats(Ports) ->
Count = length(Ports),
Sum = lists:sum(Ports),
Mean = Sum / Count,
Variance = calculate_variance_from_mean(Ports, Mean),
StdDev = math:sqrt(Variance),
#{
mean => Mean,
stddev => StdDev,
min => lists:min(Ports),
max => lists:max(Ports),
count => Count
}.
%% @private
%% @doc Calculate variance of a list.
-spec calculate_variance([number()]) -> float().
calculate_variance([]) -> 0.0;
calculate_variance([_]) -> 0.0;
calculate_variance(List) ->
Mean = lists:sum(List) / length(List),
calculate_variance_from_mean(List, Mean).
-spec calculate_variance_from_mean([number()], float()) -> float().
calculate_variance_from_mean(List, Mean) ->
SumSquaredDiffs = lists:sum([math:pow(X - Mean, 2) || X <- List]),
SumSquaredDiffs / max(1, length(List) - 1).
%% @private
%% @doc Predict ports from statistics.
-spec predict_from_stats(port_stats(), inet:port_number() | undefined, pos_integer()) ->
[inet:port_number()].
predict_from_stats(#{mean := Mean, stddev := StdDev, min := Min, max := Max}, BasePort, Count) ->
%% Generate ports within 2 standard deviations of mean
Center = case BasePort of
undefined -> round(Mean);
_ -> BasePort
end,
%% Spread predictions across the observed range
Range = max(1, Max - Min),
Step = max(1, Range div (Count + 1)),
%% Generate ports centered around the mean/base
HalfCount = Count div 2,
Ports = [clamp_port(Center + ((I - HalfCount) * Step)) || I <- lists:seq(0, Count - 1)],
%% Also include ports at +/- 1 and 2 stddev if we have room
StdDevPorts = case StdDev > 10 of
true ->
[clamp_port(round(Mean + StdDev)),
clamp_port(round(Mean - StdDev))];
false ->
[]
end,
lists:usort(Ports ++ StdDevPorts).
%% @private
%% @doc Calculate confidence for RD prediction.
-spec calculate_rd_confidence(port_stats()) -> float().
calculate_rd_confidence(#{stddev := StdDev, count := Count}) when Count >= 5 ->
%% More samples and tighter distribution = higher confidence
CountFactor = min(1.0, Count / 10),
StdDevFactor = max(0.1, 1.0 - (StdDev / 10000)),
CountFactor * StdDevFactor * 0.5; % Cap at 0.5 for RD
calculate_rd_confidence(#{count := Count}) when Count >= 3 ->
0.25;
calculate_rd_confidence(_) ->
0.1.
%% @private
%% @doc Infer allocation policy from port history.
-spec infer_allocation_policy(ets:tid(), binary()) -> {ok, allocation_policy()} | unknown.
infer_allocation_policy(Table, NodeId) ->
case ets:lookup(Table, NodeId) of
[{NodeId, #{ports := Ports, deltas := Deltas}}] when length(Ports) >= 3 ->
analyze_port_pattern(Ports, Deltas);
_ ->
unknown
end.
%% @private
%% @doc Analyze port pattern to determine allocation policy.
-spec analyze_port_pattern([inet:port_number()], [integer()]) ->
{ok, allocation_policy()} | unknown.
analyze_port_pattern(Ports, Deltas) ->
%% Check for port preservation (all same or very close)
UniqueCount = length(lists:usort(Ports)),
case UniqueCount of
1 -> {ok, pp};
_ ->
case length(Deltas) >= 2 of
true ->
%% Check for port contiguity (consistent delta)
Variance = calculate_variance(Deltas),
case Variance < 10 of % Low variance = consistent delta
true -> {ok, pc};
false -> {ok, rd}
end;
false ->
unknown
end
end.
%%%===================================================================
%%% Internal functions - Utilities
%%%===================================================================
%% @private
%% @doc Clamp port to valid range.
-spec clamp_port(integer()) -> inet:port_number().
clamp_port(Port) when Port < ?MIN_PORT -> ?MIN_PORT;
clamp_port(Port) when Port > ?MAX_PORT -> ?MAX_PORT;
clamp_port(Port) -> Port.
%% @private
%% @doc Schedule periodic cleanup.
-spec schedule_cleanup() -> reference().
schedule_cleanup() ->
erlang:send_after(?CLEANUP_INTERVAL_MS, self(), cleanup).
%% @private
%% @doc Clean up expired history entries.
-spec cleanup_expired_history(ets:tid()) -> ok.
cleanup_expired_history(Table) ->
Now = erlang:system_time(second),
MaxAge = ?HISTORY_TTL_SECONDS,
ExpiredKeys = ets:foldl(
fun({NodeId, #{updated_at := UpdatedAt}}, Acc) ->
case Now - UpdatedAt > MaxAge of
true -> [NodeId | Acc];
false -> Acc
end
end,
[],
Table
),
lists:foreach(fun(Key) -> ets:delete(Table, Key) end, ExpiredKeys),
case ExpiredKeys of
[] -> ok;
_ -> ?LOG_DEBUG("Port predictor cleanup: removed ~p expired entries",
[length(ExpiredKeys)])
end.