-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.