Current section

Files

Jump to
gaffer src gaffer_queue_runner.erl
Raw

src/gaffer_queue_runner.erl

-module(gaffer_queue_runner).
-moduledoc false.
-behaviour(gen_statem).
% API
-ignore_xref(start_link/1).
-export([start_link/1]).
-ignore_xref(poll/1).
-export([poll/1]).
-ignore_xref(reconfigure/1).
-export([reconfigure/1]).
-ignore_xref(info/1).
-export([info/1]).
% gen_statem Callbacks
-export([callback_mode/0]).
-export([init/1]).
-export([handle_event/4]).
%--- API -----------------------------------------------------------------------
-spec start_link(gaffer:queue()) -> gen_statem:start_ret().
start_link(Name) ->
gen_statem:start_link({local, proc_name(Name)}, ?MODULE, Name, []).
-spec poll(gaffer:queue()) -> ok.
poll(Name) -> call(Name, poll).
-spec reconfigure(gaffer:queue()) -> ok.
reconfigure(Name) -> call(Name, reconfigure).
-spec info(gaffer:queue()) -> non_neg_integer().
info(Name) -> call(Name, info).
%--- gen_statem Callbacks ------------------------------------------------------
callback_mode() -> handle_event_function.
init(Name) ->
Conf = gaffer_queue:conf(Name),
Data = #{name => Name, workers => #{}},
{ok, polling, Data, initial_poll(Conf) ++ poll_timeout(Conf)}.
handle_event(internal, poll, _State, #{name := Name} = Data) ->
Conf = gaffer_queue:conf(Name),
{keep_state, do_poll(Conf, Data)};
handle_event(state_timeout, poll, _State, #{name := Name} = Data) ->
Conf = gaffer_queue:conf(Name),
{keep_state, do_poll(Conf, Data), poll_timeout(Conf)};
handle_event({call, From}, poll, _State, #{name := Name} = Data) ->
Conf = gaffer_queue:conf(Name),
{keep_state, do_poll(Conf, Data), [{reply, From, ok} | poll_timeout(Conf)]};
handle_event({call, From}, reconfigure, _State, #{name := Name}) ->
{keep_state_and_data, [
{reply, From, ok} | poll_timeout(gaffer_queue:conf(Name))
]};
handle_event({call, From}, info, _State, #{workers := Workers}) ->
{keep_state_and_data, [{reply, From, map_size(Workers)}]};
handle_event(
info,
{'DOWN', _Ref, process, Pid, Reason},
_State,
#{name := Name, workers := Workers} = Data
) ->
case maps:take(Pid, Workers) of
{OriginalJob, Workers1} ->
_ = handle_worker_result(Reason, OriginalJob, Name),
{keep_state, Data#{workers := Workers1}};
error ->
{keep_state, Data}
end.
%--- Internal ------------------------------------------------------------------
call(Name, Msg) -> gen_statem:call(proc_name(Name), Msg).
initial_poll(#{poll_interval := infinity}) -> [];
initial_poll(#{}) -> [{next_event, internal, poll}].
poll_timeout(#{poll_interval := infinity}) -> [];
poll_timeout(#{poll_interval := Interval}) -> [{state_timeout, Interval, poll}].
do_poll(
#{max_workers := MaxWorkers, worker := Worker},
#{name := Name, workers := Workers} = Data
) ->
case poll_limit(MaxWorkers, map_size(Workers)) of
0 ->
Data;
Limit ->
Jobs = gaffer_queue:claim_jobs(Name, #{
queue => Name, limit => Limit
}),
NewWorkers = spawn_workers(Worker, Jobs, Workers),
Data#{workers := NewWorkers}
end.
poll_limit(infinity, _Active) -> infinity;
poll_limit(Max, Active) -> max(0, Max - Active).
% elp:ignore W0048 - spawn_workers calls gaffer_job:execute which is no_return
-dialyzer({no_return, [spawn_workers/3]}).
spawn_workers(_Worker, [], Workers) ->
Workers;
spawn_workers(Worker, [Job | Rest], Workers) ->
{Pid, _Ref} = spawn_monitor(fun() ->
gaffer_job:execute(Worker, Job)
end),
spawn_workers(Worker, Rest, Workers#{Pid => Job}).
handle_worker_result({gaffer_job, Event, Job}, _OriginalJob, Name) ->
gaffer_queue:write_result(Name, Event, Job);
handle_worker_result(Reason, OriginalJob, Name) ->
{Event, Job} = gaffer_job:handle_crash(OriginalJob, Reason),
gaffer_queue:write_result(Name, Event, Job).
proc_name(Name) ->
% elp:ignore W0023 - bounded by queue count, not user input
binary_to_atom(
<<"gaffer_queue_runner_", (atom_to_binary(Name))/binary>>
).