Current section
Files
Jump to
Current section
Files
src/glimit@rate_limiter.erl
-module(glimit@rate_limiter).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch, inline]).
-define(FILEPATH, "src/glimit/rate_limiter.gleam").
-export([new/3, hit/2, get_count/1, remove/2, sweep/1, set_now/2]).
-export_type([hit_error/0, state/1, message/1]).
-if(?OTP_RELEASE >= 27).
-define(MODULEDOC(Str), -moduledoc(Str)).
-define(DOC(Str), -doc(Str)).
-else.
-define(MODULEDOC(Str), -compile([])).
-define(DOC(Str), -compile([])).
-endif.
?MODULEDOC(false).
-type hit_error() :: rate_limited | unavailable.
-type state(FOD) :: {state,
fun((FOD) -> integer()),
fun((FOD) -> integer()),
gleam@dict:dict(FOD, glimit@bucket:bucket_state()),
gleam@option:option(integer()),
gleam@option:option(integer()),
gleam@erlang@process:subject(message(FOD)),
gleam@option:option(integer())}.
-type message(FOE) :: {hit,
FOE,
gleam@erlang@process:subject({ok, nil} | {error, hit_error()})} |
{get_count, gleam@erlang@process:subject(integer())} |
{remove, FOE, gleam@erlang@process:subject(nil)} |
sweep |
{sweep_sync, gleam@erlang@process:subject(nil)} |
{set_now, integer(), gleam@erlang@process:subject(nil)}.
-file("src/glimit/rate_limiter.gleam", 75).
?DOC(false).
-spec get_now(state(any())) -> integer().
get_now(State) ->
case erlang:element(8, State) of
{some, Now} ->
Now;
none ->
glimit@utils:now()
end.
-file("src/glimit/rate_limiter.gleam", 82).
?DOC(false).
-spec ensure_bucket(state(FOK), FOK) -> {ok,
{glimit@bucket:bucket_state(), state(FOK)}} |
{error, nil}.
ensure_bucket(State, Identifier) ->
case gleam_stdlib:map_get(erlang:element(4, State), Identifier) of
{ok, B} ->
{ok, {B, State}};
{error, _} ->
gleam@result:'try'(
glimit_ffi:rescue(
fun() -> (erlang:element(2, State))(Identifier) end
),
fun(Max) ->
gleam@result:'try'(
glimit_ffi:rescue(
fun() -> (erlang:element(3, State))(Identifier) end
),
fun(Rate) -> case glimit@bucket:new(Max, Rate) of
{ok, B@1} ->
Buckets = gleam@dict:insert(
erlang:element(4, State),
Identifier,
B@1
),
{ok,
{B@1,
{state,
erlang:element(2, State),
erlang:element(3, State),
Buckets,
erlang:element(5, State),
erlang:element(6, State),
erlang:element(7, State),
erlang:element(8, State)}}};
{error, _} ->
{error, nil}
end end
)
end
)
end.
-file("src/glimit/rate_limiter.gleam", 116).
?DOC(false).
-spec is_idle(
glimit@bucket:bucket_state(),
integer(),
gleam@option:option(integer())
) -> boolean().
is_idle(B, Now, Max_idle_ms) ->
case Max_idle_ms of
none ->
false;
{some, Threshold} ->
case erlang:element(5, B) of
none ->
true;
{some, Last_update} ->
(Now - Last_update) > Threshold
end
end.
-file("src/glimit/rate_limiter.gleam", 106).
?DOC(false).
-spec do_sweep(state(FOP)) -> state(FOP).
do_sweep(State) ->
Now = get_now(State),
Buckets = begin
_pipe = erlang:element(4, State),
gleam@dict:filter(
_pipe,
fun(_, B) ->
not glimit@bucket:is_full(B, Now) andalso not is_idle(
B,
Now,
erlang:element(6, State)
)
end
)
end,
{state,
erlang:element(2, State),
erlang:element(3, State),
Buckets,
erlang:element(5, State),
erlang:element(6, State),
erlang:element(7, State),
erlang:element(8, State)}.
-file("src/glimit/rate_limiter.gleam", 128).
?DOC(false).
-spec schedule_sweep(state(any())) -> nil.
schedule_sweep(State) ->
case erlang:element(5, State) of
{some, Ms} ->
_ = gleam@erlang@process:send_after(
erlang:element(7, State),
Ms,
sweep
),
nil;
none ->
nil
end.
-file("src/glimit/rate_limiter.gleam", 138).
?DOC(false).
-spec handle_message(state(FOV), message(FOV)) -> gleam@otp@actor:next(state(FOV), message(FOV)).
handle_message(State, Message) ->
case Message of
{hit, Identifier, Client} ->
case ensure_bucket(State, Identifier) of
{ok, {B, State@1}} ->
Now = get_now(State@1),
{Result, B@1} = glimit@bucket:hit(B, Now),
Buckets = gleam@dict:insert(
erlang:element(4, State@1),
Identifier,
B@1
),
State@2 = {state,
erlang:element(2, State@1),
erlang:element(3, State@1),
Buckets,
erlang:element(5, State@1),
erlang:element(6, State@1),
erlang:element(7, State@1),
erlang:element(8, State@1)},
case Result of
{ok, nil} ->
gleam@otp@actor:send(Client, {ok, nil});
{error, nil} ->
gleam@otp@actor:send(Client, {error, rate_limited})
end,
gleam@otp@actor:continue(State@2);
{error, _} ->
gleam@otp@actor:send(Client, {error, unavailable}),
gleam@otp@actor:continue(State)
end;
{get_count, Client@1} ->
gleam@otp@actor:send(Client@1, maps:size(erlang:element(4, State))),
gleam@otp@actor:continue(State);
{remove, Identifier@1, Client@2} ->
Buckets@1 = gleam@dict:delete(
erlang:element(4, State),
Identifier@1
),
State@3 = {state,
erlang:element(2, State),
erlang:element(3, State),
Buckets@1,
erlang:element(5, State),
erlang:element(6, State),
erlang:element(7, State),
erlang:element(8, State)},
gleam@otp@actor:send(Client@2, nil),
gleam@otp@actor:continue(State@3);
sweep ->
State@4 = do_sweep(State),
schedule_sweep(State@4),
gleam@otp@actor:continue(State@4);
{sweep_sync, Client@3} ->
State@5 = do_sweep(State),
gleam@otp@actor:send(Client@3, nil),
gleam@otp@actor:continue(State@5);
{set_now, Now@1, Client@4} ->
gleam@otp@actor:send(Client@4, nil),
gleam@otp@actor:continue(
{state,
erlang:element(2, State),
erlang:element(3, State),
erlang:element(4, State),
erlang:element(5, State),
erlang:element(6, State),
erlang:element(7, State),
{some, Now@1}}
)
end.
-file("src/glimit/rate_limiter.gleam", 196).
?DOC(false).
-spec new(
fun((FPC) -> integer()),
fun((FPC) -> integer()),
gleam@option:option(integer())
) -> {ok, gleam@erlang@process:subject(message(FPC))} | {error, nil}.
new(Per_second, Burst_limit, Max_idle_ms) ->
Sweep_interval_ms = {some, 10000},
_pipe@1 = gleam@otp@actor:new_with_initialiser(
1000,
fun(Self_subject) ->
State = {state,
Burst_limit,
Per_second,
maps:new(),
Sweep_interval_ms,
Max_idle_ms,
Self_subject,
none},
schedule_sweep(State),
{ok,
begin
_pipe = gleam@otp@actor:initialised(State),
gleam@otp@actor:returning(_pipe, Self_subject)
end}
end
),
_pipe@2 = gleam@otp@actor:on_message(_pipe@1, fun handle_message/2),
_pipe@3 = gleam@otp@actor:start(_pipe@2),
_pipe@4 = gleam@result:map(
_pipe@3,
fun(Started) -> erlang:element(3, Started) end
),
gleam@result:map_error(_pipe@4, fun(_) -> nil end).
-file("src/glimit/rate_limiter.gleam", 230).
?DOC(false).
-spec hit(gleam@erlang@process:subject(message(FPH)), FPH) -> {ok, nil} |
{error, hit_error()}.
hit(Rate_limiter, Identifier) ->
case glimit@utils:safe_call(
Rate_limiter,
fun(_capture) -> {hit, Identifier, _capture} end,
1000
) of
{ok, {ok, nil}} ->
{ok, nil};
{ok, {error, Err}} ->
{error, Err};
{error, nil} ->
{error, unavailable}
end.
-file("src/glimit/rate_limiter.gleam", 243).
?DOC(false).
-spec get_count(gleam@erlang@process:subject(message(any()))) -> integer().
get_count(Rate_limiter) ->
_pipe = glimit@utils:safe_call(
Rate_limiter,
fun(Field@0) -> {get_count, Field@0} end,
1000
),
gleam@result:unwrap(_pipe, 0).
-file("src/glimit/rate_limiter.gleam", 250).
?DOC(false).
-spec remove(gleam@erlang@process:subject(message(FPN)), FPN) -> {ok, nil} |
{error, nil}.
remove(Rate_limiter, Identifier) ->
glimit@utils:safe_call(
Rate_limiter,
fun(_capture) -> {remove, Identifier, _capture} end,
1000
).
-file("src/glimit/rate_limiter.gleam", 260).
?DOC(false).
-spec sweep(gleam@erlang@process:subject(message(any()))) -> {ok, nil} |
{error, nil}.
sweep(Rate_limiter) ->
glimit@utils:safe_call(
Rate_limiter,
fun(Field@0) -> {sweep_sync, Field@0} end,
1000
).
-file("src/glimit/rate_limiter.gleam", 267).
?DOC(false).
-spec set_now(gleam@erlang@process:subject(message(any())), integer()) -> nil.
set_now(Rate_limiter, Now) ->
_ = glimit@utils:safe_call(
Rate_limiter,
fun(_capture) -> {set_now, Now, _capture} end,
1000
),
nil.