-
-
Notifications
You must be signed in to change notification settings - Fork 453
Expand file tree
/
Copy pathhackney_load_regulation.erl
More file actions
149 lines (135 loc) · 4.51 KB
/
Copy pathhackney_load_regulation.erl
File metadata and controls
149 lines (135 loc) · 4.51 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
%%% -*- erlang -*-
%%%
%%% This file is part of hackney released under the Apache 2 license.
%%% See the NOTICE for more information.
%%%
%%% Copyright (c) 2024-2026 Benoit Chesneau
%%%
%%% @doc Per-host connection load regulation using ETS counting semaphore.
%%%
%%% This module provides per-host connection limits using an atomic
%%% counting semaphore pattern. It tracks the number of active connections
%%% per {Host, Port} and blocks new requests when the limit is reached.
%%%
%%% Usage:
%%% ```
%%% case hackney_load_regulation:acquire(Host, Port, MaxPerHost, Timeout) of
%%% ok ->
%%% try
%%% %% Do work with connection
%%% after
%%% hackney_load_regulation:release(Host, Port)
%%% end;
%%% {error, timeout} ->
%%% {error, checkout_timeout}
%%% end.
%%% '''
-module(hackney_load_regulation).
%% API
-export([
init/0,
acquire/4,
release/2,
current/2,
reset/2
]).
-define(TABLE, hackney_host_limits).
-define(BACKOFF_MS, 10).
%%====================================================================
%% API
%%====================================================================
%% @doc Initialize the load regulation ETS table.
%% Should be called once during application startup.
-spec init() -> ok.
init() ->
case ets:whereis(?TABLE) of
undefined ->
?TABLE = ets:new(?TABLE, [
public,
set,
named_table,
{write_concurrency, true},
{read_concurrency, true}
]),
ok;
_Tid ->
ok
end.
%% @doc Acquire a slot for the given host.
%% Blocks with exponential backoff until a slot is available or timeout.
%% Returns ok if slot acquired, {error, timeout} otherwise.
-spec acquire(Host :: string() | binary(), Port :: inet:port_number(),
MaxPerHost :: pos_integer(), Timeout :: timeout()) ->
ok | {error, timeout}.
acquire(Host, Port, MaxPerHost, Timeout) ->
Key = normalize_key(Host, Port),
Deadline = deadline(Timeout),
acquire_loop(Key, MaxPerHost, Deadline).
%% @doc Release a slot for the given host.
%% Should always be called after acquire, typically in an after block.
-spec release(Host :: string() | binary(), Port :: inet:port_number()) -> ok.
release(Host, Port) ->
Key = normalize_key(Host, Port),
try
_ = ets:update_counter(?TABLE, Key, {2, -1, 0, 0}),
ok
catch
error:badarg ->
%% Key doesn't exist, nothing to release
ok
end.
%% @doc Get the current number of active connections for a host.
-spec current(Host :: string() | binary(), Port :: inet:port_number()) ->
non_neg_integer().
current(Host, Port) ->
Key = normalize_key(Host, Port),
case ets:lookup(?TABLE, Key) of
[{_, Count}] -> max(0, Count);
[] -> 0
end.
%% @doc Reset the counter for a host (for testing).
-spec reset(Host :: string() | binary(), Port :: inet:port_number()) -> ok.
reset(Host, Port) ->
Key = normalize_key(Host, Port),
ets:delete(?TABLE, Key),
ok.
%%====================================================================
%% Internal functions
%%====================================================================
%% @private Normalize host to lowercase binary for consistent keys.
normalize_key(Host, Port) when is_list(Host) ->
normalize_key(list_to_binary(Host), Port);
normalize_key(Host, Port) when is_binary(Host) ->
{string:lowercase(Host), Port}.
%% @private Calculate deadline from timeout.
deadline(infinity) ->
infinity;
deadline(Timeout) when is_integer(Timeout), Timeout >= 0 ->
erlang:monotonic_time(millisecond) + Timeout.
%% @private Check if deadline has passed.
check_deadline(infinity) ->
ok;
check_deadline(Deadline) ->
case erlang:monotonic_time(millisecond) < Deadline of
true -> ok;
false -> timeout
end.
%% @private Main acquire loop with backoff.
acquire_loop(Key, Max, Deadline) ->
%% Atomically increment counter, creating entry if needed
Count = ets:update_counter(?TABLE, Key, {2, 1}, {Key, 0}),
case Count =< Max of
true ->
%% Got a slot
ok;
false ->
%% Over limit - decrement back and retry
_ = ets:update_counter(?TABLE, Key, {2, -1}),
case check_deadline(Deadline) of
ok ->
timer:sleep(?BACKOFF_MS),
acquire_loop(Key, Max, Deadline);
timeout ->
{error, timeout}
end
end.