Current section

Files

Jump to
telega src telega@internal@request_queue.erl
Raw

src/telega@internal@request_queue.erl

-module(telega@internal@request_queue).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch, inline]).
-define(FILEPATH, "src/telega/internal/request_queue.gleam").
-export([default_config/0, start/1, execute_with_rule/4, execute/2, shutdown/1, total_length/1, is_overheated/1]).
-export_type([request_queue/0, rule/0, queue_config/0, queued_request/0, message/0, rule_state/0, state/0]).
-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).
-opaque request_queue() :: {request_queue,
gleam@erlang@process:subject(message())}.
-type rule() :: {rule, binary(), integer(), integer(), integer()}.
-type queue_config() :: {queue_config,
list(rule()),
gleam@option:option(integer()),
gleam@option:option(integer()),
integer(),
integer()}.
-type queued_request() :: {queued_request,
binary(),
binary(),
fun(() -> {ok, gleam@http@response:response(binary())} |
{error, telega@error:telega_error()}),
gleam@erlang@process:subject({ok,
gleam@http@response:response(binary())} |
{error, telega@error:telega_error()}),
integer()}.
-type message() :: {execute, queued_request()} |
process_queue |
{request_completed, binary(), binary()} |
{request_failed, queued_request(), telega@error:telega_error(), boolean()} |
{retry_request, queued_request()} |
{get_total_length, gleam@erlang@process:subject(integer())} |
{is_overheated, gleam@erlang@process:subject(boolean())} |
shutdown.
-type rule_state() :: {rule_state,
rule(),
integer(),
integer(),
list(queued_request())}.
-type state() :: {state,
queue_config(),
gleam@dict:dict(binary(), rule_state()),
integer(),
integer(),
gleam@dict:dict(binary(), binary()),
gleam@erlang@process:subject(message())}.
-file("src/telega/internal/request_queue.gleam", 49).
?DOC(false).
-spec default_config() -> queue_config().
default_config() ->
{queue_config,
[{rule, <<"default"/utf8>>, 30, 1000, 5}],
{some, 30},
{some, 100},
1000,
3}.
-file("src/telega/internal/request_queue.gleam", 294).
?DOC(false).
-spec emit_queue_depth(rule(), integer()) -> nil.
emit_queue_depth(Rule, Depth) ->
telega@telemetry:execute(
[<<"telega"/utf8>>, <<"request_queue"/utf8>>, <<"depth"/utf8>>],
[{<<"depth"/utf8>>, Depth}],
[{<<"rule_id"/utf8>>, {string_value, erlang:element(2, Rule)}},
{<<"priority"/utf8>>, {int_value, erlang:element(5, Rule)}}]
).
-file("src/telega/internal/request_queue.gleam", 301).
?DOC(false).
-spec add_to_queue(state(), queued_request()) -> state().
add_to_queue(State, Request) ->
case gleam_stdlib:map_get(
erlang:element(3, State),
erlang:element(3, Request)
) of
{ok, Rule_state} ->
New_queue = lists:append(erlang:element(5, Rule_state), [Request]),
New_rule_state = {rule_state,
erlang:element(2, Rule_state),
erlang:element(3, Rule_state),
erlang:element(4, Rule_state),
New_queue},
New_rule_states = gleam@dict:insert(
erlang:element(3, State),
erlang:element(3, Request),
New_rule_state
),
emit_queue_depth(
erlang:element(2, Rule_state),
erlang:length(New_queue)
),
{state,
erlang:element(2, State),
New_rule_states,
erlang:element(4, State),
erlang:element(5, State),
erlang:element(6, State),
erlang:element(7, State)};
{error, _} ->
case gleam_stdlib:map_get(
erlang:element(3, State),
<<"default"/utf8>>
) of
{ok, Rule_state@1} ->
Updated_request = {queued_request,
erlang:element(2, Request),
<<"default"/utf8>>,
erlang:element(4, Request),
erlang:element(5, Request),
erlang:element(6, Request)},
New_queue@1 = lists:append(
erlang:element(5, Rule_state@1),
[Updated_request]
),
New_rule_state@1 = {rule_state,
erlang:element(2, Rule_state@1),
erlang:element(3, Rule_state@1),
erlang:element(4, Rule_state@1),
New_queue@1},
New_rule_states@1 = gleam@dict:insert(
erlang:element(3, State),
<<"default"/utf8>>,
New_rule_state@1
),
emit_queue_depth(
erlang:element(2, Rule_state@1),
erlang:length(New_queue@1)
),
{state,
erlang:element(2, State),
New_rule_states@1,
erlang:element(4, State),
erlang:element(5, State),
erlang:element(6, State),
erlang:element(7, State)};
{error, _} ->
gleam@erlang@process:send(
erlang:element(5, Request),
{error, {fetch_error, <<"Invalid rule ID"/utf8>>}}
),
State
end
end.
-file("src/telega/internal/request_queue.gleam", 425).
?DOC(false).
-spec execute_request(
queued_request(),
gleam@erlang@process:subject(message()),
integer()
) -> nil.
execute_request(Request, Self, Max_retries) ->
Result = (erlang:element(4, Request))(),
case Result of
{ok, Value} ->
gleam@erlang@process:send(
Self,
{request_completed,
erlang:element(2, Request),
erlang:element(3, Request)}
),
gleam@erlang@process:send(erlang:element(5, Request), {ok, Value});
{error, Error} ->
Should_retry = erlang:element(6, Request) < Max_retries,
case Should_retry of
true ->
gleam@erlang@process:send(
Self,
{request_failed, Request, Error, true}
);
false ->
gleam@erlang@process:send(
Self,
{request_completed,
erlang:element(2, Request),
erlang:element(3, Request)}
),
gleam@erlang@process:send(
erlang:element(5, Request),
{error, Error}
)
end
end,
nil.
-file("src/telega/internal/request_queue.gleam", 392).
?DOC(false).
-spec can_process(state(), rule_state(), integer()) -> boolean().
can_process(State, Rule_state, _) ->
Rule_ok = erlang:element(3, Rule_state) < erlang:element(
3,
erlang:element(2, Rule_state)
),
Overall_ok = case erlang:element(3, erlang:element(2, State)) of
{some, Limit} ->
erlang:element(4, State) < Limit;
none ->
true
end,
Concurrent_ok = case erlang:element(4, erlang:element(2, State)) of
{some, Limit@1} ->
maps:size(erlang:element(6, State)) < Limit@1;
none ->
true
end,
(Rule_ok andalso Overall_ok) andalso Concurrent_ok.
-file("src/telega/internal/request_queue.gleam", 355).
?DOC(false).
-spec process_rule_queue(state(), binary(), rule_state(), integer()) -> state().
process_rule_queue(State, Rule_id, Rule_state, Now) ->
case erlang:element(5, Rule_state) of
[] ->
State;
[Request | Rest] ->
case can_process(State, Rule_state, Now) of
true ->
execute_request(
Request,
erlang:element(7, State),
erlang:element(6, erlang:element(2, State))
),
emit_queue_depth(
erlang:element(2, Rule_state),
erlang:length(Rest)
),
New_rule_state = {rule_state,
erlang:element(2, Rule_state),
erlang:element(3, Rule_state) + 1,
erlang:element(4, Rule_state),
Rest},
New_rule_states = gleam@dict:insert(
erlang:element(3, State),
Rule_id,
New_rule_state
),
New_in_flight = gleam@dict:insert(
erlang:element(6, State),
erlang:element(2, Request),
Rule_id
),
{state,
erlang:element(2, State),
New_rule_states,
erlang:element(4, State) + 1,
erlang:element(5, State),
New_in_flight,
erlang:element(7, State)};
false ->
State
end
end.
-file("src/telega/internal/request_queue.gleam", 408).
?DOC(false).
-spec reset_windows(state(), integer()) -> state().
reset_windows(State, Now) ->
State@1 = case (Now - erlang:element(5, State)) > 1000 of
true ->
{state,
erlang:element(2, State),
erlang:element(3, State),
0,
Now,
erlang:element(6, State),
erlang:element(7, State)};
false ->
State
end,
New_rule_states = gleam@dict:map_values(
erlang:element(3, State@1),
fun(_, Rule_state) ->
case (Now - erlang:element(4, Rule_state)) > erlang:element(
4,
erlang:element(2, Rule_state)
) of
true ->
{rule_state,
erlang:element(2, Rule_state),
0,
Now,
erlang:element(5, Rule_state)};
false ->
Rule_state
end
end
),
{state,
erlang:element(2, State@1),
New_rule_states,
erlang:element(4, State@1),
erlang:element(5, State@1),
erlang:element(6, State@1),
erlang:element(7, State@1)}.
-file("src/telega/internal/request_queue.gleam", 336).
?DOC(false).
-spec process_all_queues(state()) -> state().
process_all_queues(State) ->
Now = telega@internal@utils:current_time_ms(),
State@1 = reset_windows(State, Now),
Sorted_rules = begin
_pipe = maps:to_list(erlang:element(3, State@1)),
gleam@list:sort(
_pipe,
fun(A, B) ->
{_, Rule_state_a} = A,
{_, Rule_state_b} = B,
gleam@int:compare(
erlang:element(5, erlang:element(2, Rule_state_a)),
erlang:element(5, erlang:element(2, Rule_state_b))
)
end
)
end,
gleam@list:fold(
Sorted_rules,
State@1,
fun(State@2, Rule_entry) ->
{Rule_id, Rule_state} = Rule_entry,
process_rule_queue(State@2, Rule_id, Rule_state, Now)
end
).
-file("src/telega/internal/request_queue.gleam", 210).
?DOC(false).
-spec handle_message(state(), message()) -> gleam@otp@actor:next(state(), message()).
handle_message(State, Message) ->
case Message of
{execute, Request} ->
New_state = add_to_queue(State, Request),
gleam@erlang@process:send(
erlang:element(7, New_state),
process_queue
),
gleam@otp@actor:continue(New_state);
process_queue ->
New_state@1 = process_all_queues(State),
gleam@erlang@process:send_after(
erlang:element(7, State),
100,
process_queue
),
gleam@otp@actor:continue(New_state@1);
{request_completed, Id, _} ->
New_in_flight = gleam@dict:delete(erlang:element(6, State), Id),
New_state@2 = {state,
erlang:element(2, State),
erlang:element(3, State),
erlang:element(4, State),
erlang:element(5, State),
New_in_flight,
erlang:element(7, State)},
gleam@erlang@process:send(
erlang:element(7, New_state@2),
process_queue
),
gleam@otp@actor:continue(New_state@2);
{request_failed, Request@1, Error, Should_retry} ->
New_in_flight@1 = gleam@dict:delete(
erlang:element(6, State),
erlang:element(2, Request@1)
),
Mut_state = {state,
erlang:element(2, State),
erlang:element(3, State),
erlang:element(4, State),
erlang:element(5, State),
New_in_flight@1,
erlang:element(7, State)},
case Should_retry of
true ->
Retry_request = {queued_request,
erlang:element(2, Request@1),
erlang:element(3, Request@1),
erlang:element(4, Request@1),
erlang:element(5, Request@1),
erlang:element(6, Request@1) + 1},
gleam@erlang@process:send_after(
erlang:element(7, State),
erlang:element(5, erlang:element(2, State)),
{retry_request, Retry_request}
),
gleam@otp@actor:continue(Mut_state);
false ->
gleam@erlang@process:send(
erlang:element(5, Request@1),
{error, Error}
),
gleam@otp@actor:continue(Mut_state)
end;
{retry_request, Request@2} ->
New_state@3 = add_to_queue(State, Request@2),
gleam@erlang@process:send(
erlang:element(7, New_state@3),
process_queue
),
gleam@otp@actor:continue(New_state@3);
{get_total_length, Reply_to} ->
Total = gleam@dict:fold(
erlang:element(3, State),
0,
fun(Acc, _, Rule_state) ->
Acc + erlang:length(erlang:element(5, Rule_state))
end
),
gleam@erlang@process:send(Reply_to, Total),
gleam@otp@actor:continue(State);
{is_overheated, Reply_to@1} ->
Overheated = begin
_pipe = maps:to_list(erlang:element(3, State)),
gleam@list:any(
_pipe,
fun(Pair) ->
{_, Rule_state@1} = Pair,
erlang:element(3, Rule_state@1) >= erlang:element(
3,
erlang:element(2, Rule_state@1)
)
end
)
end,
gleam@erlang@process:send(Reply_to@1, Overheated),
gleam@otp@actor:continue(State);
shutdown ->
gleam@otp@actor:stop()
end.
-file("src/telega/internal/request_queue.gleam", 116).
?DOC(false).
-spec start(queue_config()) -> {ok, request_queue()} |
{error, gleam@otp@actor:start_error()}.
start(Config) ->
gleam@result:'try'(
begin
_pipe@2 = gleam@otp@actor:new_with_initialiser(
1000,
fun(Self) ->
Rule_states = gleam@list:fold(
erlang:element(2, Config),
maps:new(),
fun(Acc, Rule) ->
gleam@dict:insert(
Acc,
erlang:element(2, Rule),
{rule_state, Rule, 0, 0, []}
)
end
),
Initial_state = {state,
Config,
Rule_states,
0,
0,
maps:new(),
Self},
gleam@erlang@process:send_after(Self, 100, process_queue),
_pipe = gleam@otp@actor:initialised(Initial_state),
_pipe@1 = gleam@otp@actor:returning(_pipe, Self),
{ok, _pipe@1}
end
),
_pipe@3 = gleam@otp@actor:on_message(_pipe@2, fun handle_message/2),
_pipe@4 = gleam@otp@actor:start(_pipe@3),
gleam@result:map_error(_pipe@4, fun(_) -> init_timeout end)
end,
fun(Started) -> {ok, {request_queue, erlang:element(3, Started)}} end
).
-file("src/telega/internal/request_queue.gleam", 153).
?DOC(false).
-spec execute_with_rule(
request_queue(),
binary(),
binary(),
fun(() -> {ok, gleam@http@response:response(binary())} |
{error, telega@error:telega_error()})
) -> {ok, gleam@http@response:response(binary())} |
{error, telega@error:telega_error()}.
execute_with_rule(Queue, Request_id, Rule_id, Execute) ->
Reply_subject = gleam@erlang@process:new_subject(),
Request = {queued_request, Request_id, Rule_id, Execute, Reply_subject, 0},
gleam@erlang@process:send(erlang:element(2, Queue), {execute, Request}),
gleam_erlang_ffi:'receive'(Reply_subject).
-file("src/telega/internal/request_queue.gleam", 176).
?DOC(false).
-spec execute(
request_queue(),
fun(() -> {ok, gleam@http@response:response(binary())} |
{error, telega@error:telega_error()})
) -> {ok, gleam@http@response:response(binary())} |
{error, telega@error:telega_error()}.
execute(Queue, Execute) ->
execute_with_rule(
Queue,
telega@internal@utils:random_string(32),
<<"default"/utf8>>,
Execute
).
-file("src/telega/internal/request_queue.gleam", 184).
?DOC(false).
-spec shutdown(request_queue()) -> nil.
shutdown(Queue) ->
gleam@erlang@process:send(erlang:element(2, Queue), shutdown).
-file("src/telega/internal/request_queue.gleam", 189).
?DOC(false).
-spec total_length(request_queue()) -> integer().
total_length(Queue) ->
Reply_subject = gleam@erlang@process:new_subject(),
gleam@erlang@process:send(
erlang:element(2, Queue),
{get_total_length, Reply_subject}
),
case gleam@erlang@process:'receive'(Reply_subject, 1000) of
{ok, Length} ->
Length;
{error, _} ->
0
end.
-file("src/telega/internal/request_queue.gleam", 200).
?DOC(false).
-spec is_overheated(request_queue()) -> boolean().
is_overheated(Queue) ->
Reply_subject = gleam@erlang@process:new_subject(),
gleam@erlang@process:send(
erlang:element(2, Queue),
{is_overheated, Reply_subject}
),
case gleam@erlang@process:'receive'(Reply_subject, 1000) of
{ok, Overheated} ->
Overheated;
{error, _} ->
false
end.