Current section

Files

Jump to
glimit src glimit@registry.erl
Raw

src/glimit@registry.erl

-module(glimit@registry).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch]).
-export([get_or_create/2, get_all/1, remove/2, sweep/2, new/2]).
-export_type([state/1, message/1]).
-type state(HPL) :: {state,
fun((HPL) -> integer()),
fun((HPL) -> integer()),
gleam@dict:dict(HPL, gleam@erlang@process:subject(glimit@rate_limiter:message()))}.
-type message(HPM) :: {get_or_create,
HPM,
gleam@erlang@process:subject({ok,
gleam@erlang@process:subject(glimit@rate_limiter:message())} |
{error, nil})} |
{get_all,
gleam@erlang@process:subject(list({HPM,
gleam@erlang@process:subject(glimit@rate_limiter:message())}))} |
{remove, HPM, gleam@erlang@process:subject(nil)}.
-spec handle_get_or_create(HPR, state(HPR)) -> {ok,
gleam@erlang@process:subject(glimit@rate_limiter:message())} |
{error, nil}.
handle_get_or_create(Identifier, State) ->
case begin
_pipe = erlang:element(4, State),
gleam@dict:get(_pipe, Identifier)
end of
{ok, Rate_limiter} ->
{ok, Rate_limiter};
{error, _} ->
gleam@result:'try'(
glimit@rate_limiter:new(
(erlang:element(2, State))(Identifier),
(erlang:element(3, State))(Identifier)
),
fun(Rate_limiter@1) -> {ok, Rate_limiter@1} end
)
end.
-spec handle_message(message(HPW), state(HPW)) -> gleam@otp@actor:next(message(HPW), state(HPW)).
handle_message(Message, State) ->
case Message of
{get_or_create, Identifier, Client} ->
case handle_get_or_create(Identifier, State) of
{ok, Rate_limiter} ->
Registry = begin
_pipe = erlang:element(4, State),
gleam@dict:insert(_pipe, Identifier, Rate_limiter)
end,
State@1 = erlang:setelement(4, State, Registry),
gleam@otp@actor:send(Client, {ok, Rate_limiter}),
gleam@otp@actor:continue(State@1);
{error, _} ->
gleam@otp@actor:send(Client, {error, nil}),
gleam@otp@actor:continue(State)
end;
{get_all, Client@1} ->
Rate_limiters = begin
_pipe@1 = erlang:element(4, State),
maps:to_list(_pipe@1)
end,
gleam@otp@actor:send(Client@1, Rate_limiters),
gleam@otp@actor:continue(State);
{remove, Identifier@1, Client@2} ->
Registry@1 = begin
_pipe@2 = erlang:element(4, State),
gleam@dict:delete(_pipe@2, Identifier@1)
end,
State@2 = erlang:setelement(4, State, Registry@1),
gleam@otp@actor:send(Client@2, nil),
gleam@otp@actor:continue(State@2)
end.
-spec get_or_create(gleam@erlang@process:subject(message(HQH)), HQH) -> {ok,
gleam@erlang@process:subject(glimit@rate_limiter:message())} |
{error, nil}.
get_or_create(Registry, Identifier) ->
gleam@otp@actor:call(
Registry,
fun(_capture) -> {get_or_create, Identifier, _capture} end,
10
).
-spec get_all(gleam@erlang@process:subject(message(HQM))) -> list({HQM,
gleam@erlang@process:subject(glimit@rate_limiter:message())}).
get_all(Registry) ->
gleam@otp@actor:call(Registry, fun(Field@0) -> {get_all, Field@0} end, 10).
-spec remove(gleam@erlang@process:subject(message(HQQ)), HQQ) -> {ok, nil} |
{error, nil}.
remove(Registry, Identifier) ->
gleam@otp@actor:call(
Registry,
fun(_capture) -> {remove, Identifier, _capture} end,
10
),
{ok, nil}.
-spec sweep(
gleam@erlang@process:subject(message(any())),
gleam@option:option(integer())
) -> nil.
sweep(Registry, Interval_secs) ->
case Interval_secs of
{some, I} ->
gleam_erlang_ffi:sleep(I * 1000);
none ->
nil
end,
_pipe = get_all(Registry),
_pipe@2 = gleam@list:filter(
_pipe,
fun(Pair) ->
{_, Rate_limiter} = Pair,
_pipe@1 = Rate_limiter,
glimit@rate_limiter:has_full_bucket(_pipe@1)
end
),
gleam@list:map(
_pipe@2,
fun(Pair@1) ->
{Identifier, Rate_limiter@1} = Pair@1,
_ = remove(Registry, Identifier),
_pipe@3 = Rate_limiter@1,
glimit@rate_limiter:shutdown(_pipe@3),
Identifier
end
),
case Interval_secs of
{some, _} ->
sweep(Registry, Interval_secs);
none ->
nil
end.
-spec new(fun((HQD) -> integer()), fun((HQD) -> integer())) -> {ok,
gleam@erlang@process:subject(message(HQD))} |
{error, nil}.
new(Per_second, Burst_limit) ->
State = {state, Burst_limit, Per_second, gleam@dict:new()},
gleam@result:'try'(
begin
_pipe = gleam@otp@actor:start(State, fun handle_message/2),
gleam@result:nil_error(_pipe)
end,
fun(Registry) ->
gleam@otp@task:async(fun() -> sweep(Registry, {some, 10}) end),
{ok, Registry}
end
).