Current section
Files
Jump to
Current section
Files
src/wpool_queue_manager.erl
% This file is licensed to you under the Apache License,
% Version 2.0 (the "License"); you may not use this file
% except in compliance with the License. You may obtain
% a copy of the License at
%
% https://www.apache.org/licenses/LICENSE-2.0
%
% Unless required by applicable law or agreed to in writing,
% software distributed under the License is distributed on an
% "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
% KIND, either express or implied. See the License for the
% specific language governing permissions and limitations
% under the License.
%%% @private
-module(wpool_queue_manager).
-behaviour(gen_server).
%% api
-export([start_link/2, start_link/3]).
-export([call_available_worker/3, cast_to_available_worker/2, new_worker/2, worker_dead/2,
send_request_available_worker/3, worker_ready/2, worker_busy/2, pending_task_count/1]).
%% gen_server callbacks
-export([init/1, handle_call/3, handle_cast/2, handle_info/2]).
-record(state,
{wpool :: wpool:name(),
clients :: queue:queue({cast | {pid(), _}, term()}),
workers :: gb_sets:set(atom()),
monitors :: #{atom() := monitored_from()},
queue_type :: queue_type()}).
-opaque state() :: #state{}.
-export_type([state/0]).
-type from() :: {pid(), gen_server:reply_tag()}.
-export_type([from/0]).
-type monitored_from() :: {reference(), from()}.
-type options() :: [{option(), term()}].
-export_type([options/0]).
-type option() :: queue_type.
-type args() :: [{arg(), term()}].
-export_type([args/0]).
-type arg() :: option() | pool.
-type queue_mgr() :: atom().
-type queue_type() :: fifo | lifo.
-type worker_event() :: new_worker | worker_dead | worker_busy | worker_ready.
-export_type([worker_event/0]).
-type call_request() :: {available_worker, infinity | pos_integer()} | pending_task_count.
-export_type([call_request/0]).
-export_type([queue_mgr/0, queue_type/0]).
%%%===================================================================
%%% API
%%%===================================================================
-spec start_link(wpool:name(), queue_mgr()) ->
{ok, pid()} | {error, {already_started, pid()} | term()}.
start_link(WPool, Name) ->
start_link(WPool, Name, []).
-spec start_link(wpool:name(), queue_mgr(), options()) ->
{ok, pid()} | {error, {already_started, pid()} | term()}.
start_link(WPool, Name, Options) ->
gen_server:start_link({local, Name}, ?MODULE, [{pool, WPool} | Options], []).
%% @doc returns the first available worker in the pool
-spec call_available_worker(queue_mgr(), any(), timeout()) -> noproc | timeout | any().
call_available_worker(QueueManager, Call, Timeout) ->
case get_available_worker(QueueManager, Call, Timeout) of
{ok, TimeLeft, Worker} when TimeLeft > 0 ->
wpool_process:call(Worker, Call, TimeLeft);
{ok, _, Worker} ->
worker_ready(QueueManager, Worker),
timeout;
Other ->
Other
end.
%% @doc Casts a message to the first available worker.
%% Since we can wait forever for a wpool:cast to be delivered
%% but we don't want the caller to be blocked, this function
%% just forwards the cast when it gets the worker
-spec cast_to_available_worker(queue_mgr(), term()) -> ok.
cast_to_available_worker(QueueManager, Cast) ->
gen_server:cast(QueueManager, {cast_to_available_worker, Cast}).
%% @doc returns the first available worker in the pool
-spec send_request_available_worker(queue_mgr(), any(), timeout()) ->
noproc | timeout | gen_server:request_id().
send_request_available_worker(QueueManager, Call, Timeout) ->
case get_available_worker(QueueManager, Call, Timeout) of
{ok, _TimeLeft, Worker} ->
wpool_process:send_request(Worker, Call);
Other ->
Other
end.
%% @doc Mark a brand new worker as available
-spec new_worker(queue_mgr(), atom()) -> ok.
new_worker(QueueManager, Worker) ->
gen_server:cast(QueueManager, {new_worker, Worker}).
%% @doc Mark a worker as available
-spec worker_ready(queue_mgr(), atom()) -> ok.
worker_ready(QueueManager, Worker) ->
gen_server:cast(QueueManager, {worker_ready, Worker}).
%% @doc Mark a worker as no longer available
-spec worker_busy(queue_mgr(), atom()) -> ok.
worker_busy(QueueManager, Worker) ->
gen_server:cast(QueueManager, {worker_busy, Worker}).
%% @doc Decrement the total number of workers
-spec worker_dead(queue_mgr(), atom()) -> ok.
worker_dead(QueueManager, Worker) ->
gen_server:cast(QueueManager, {worker_dead, Worker}).
%% @doc Retrieves the number of pending tasks (used for stats)
%% @see wpool_pool:stats/1
-spec pending_task_count(queue_mgr()) -> non_neg_integer().
pending_task_count(QueueManager) ->
gen_server:call(QueueManager, pending_task_count).
%%%===================================================================
%%% gen_server callbacks
%%%===================================================================
-spec init(args()) -> {ok, state()}.
init(Args) ->
WPool = proplists:get_value(pool, Args),
QueueType = proplists:get_value(queue_type, Args),
put(pending_tasks, 0),
{ok,
#state{wpool = WPool,
clients = queue:new(),
workers = gb_sets:new(),
monitors = #{},
queue_type = QueueType}}.
-spec handle_cast({worker_event(), atom()}, state()) -> {noreply, state()}.
handle_cast({new_worker, Worker}, State) ->
handle_cast({worker_ready, Worker}, State);
handle_cast({worker_dead, Worker}, #state{workers = Workers} = State) ->
NewWorkers = gb_sets:delete_any(Worker, Workers),
{noreply, State#state{workers = NewWorkers}};
handle_cast({worker_busy, Worker}, #state{workers = Workers} = State) ->
{noreply, State#state{workers = gb_sets:delete_any(Worker, Workers)}};
handle_cast({worker_ready, Worker}, State0) ->
#state{workers = Workers,
clients = Clients,
monitors = Mons,
queue_type = QueueType} =
State0,
State =
case Mons of
#{Worker := {Ref, _Client}} ->
demonitor(Ref, [flush]),
State0#state{monitors = maps:remove(Worker, Mons)};
_ ->
State0
end,
case queue_out(Clients, QueueType) of
{empty, _Clients} ->
{noreply, State#state{workers = gb_sets:add(Worker, Workers)}};
{{value, {cast, Cast}}, NewClients} ->
dec_pending_tasks(),
ok = wpool_process:cast(Worker, Cast),
{noreply, State#state{clients = NewClients}};
{{value, {Client = {ClientPid, _}, ExpiresAt}}, NewClients} ->
dec_pending_tasks(),
NewState = State#state{clients = NewClients},
case is_process_alive(ClientPid) andalso is_expired(ExpiresAt) of
true ->
MonitorState = monitor_worker(Worker, Client, NewState),
gen_server:reply(Client, {ok, Worker}),
{noreply, MonitorState};
false ->
handle_cast({worker_ready, Worker}, NewState)
end
end;
handle_cast({cast_to_available_worker, Cast}, State) ->
#state{workers = Workers, clients = Clients} = State,
case gb_sets:is_empty(Workers) of
true ->
inc_pending_tasks(),
{noreply, State#state{clients = queue:in({cast, Cast}, Clients)}};
false ->
{Worker, NewWorkers} = gb_sets:take_smallest(Workers),
ok = wpool_process:cast(Worker, Cast),
{noreply, State#state{workers = NewWorkers}}
end.
-spec handle_call(call_request(), from(), state()) ->
{reply, {ok, atom()}, state()} | {noreply, state()}.
handle_call({available_worker, ExpiresAt}, {ClientPid, _Ref} = Client, State) ->
#state{workers = Workers, clients = Clients} = State,
case gb_sets:is_empty(Workers) of
true ->
inc_pending_tasks(),
{noreply, State#state{clients = queue:in({Client, ExpiresAt}, Clients)}};
false ->
{Worker, NewWorkers} = gb_sets:take_smallest(Workers),
%NOTE: It could've been a while since this call was made, so we check
case erlang:is_process_alive(ClientPid) andalso is_expired(ExpiresAt) of
true ->
NewState = monitor_worker(Worker, Client, State#state{workers = NewWorkers}),
{reply, {ok, Worker}, NewState};
false ->
{noreply, State}
end
end;
handle_call(pending_task_count, _From, State) ->
{reply, get(pending_tasks), State}.
-spec handle_info(any(), state()) -> {noreply, state()}.
handle_info({'DOWN', Ref, Type, {Worker, _Node}, Exit}, State) ->
handle_info({'DOWN', Ref, Type, Worker, Exit}, State);
handle_info({'DOWN', _, _, Worker, Exit}, #state{monitors = Mons} = State) ->
case Mons of
#{Worker := {_Ref, Client}} ->
gen_server:reply(Client, {'EXIT', Worker, Exit}),
{noreply, State#state{monitors = maps:remove(Worker, Mons)}};
_ ->
{noreply, State}
end;
handle_info(_Info, State) ->
{noreply, State}.
%%%===================================================================
%%% private
%%%===================================================================
-spec get_available_worker(queue_mgr(), any(), timeout()) ->
noproc | timeout | {ok, timeout(), any()}.
get_available_worker(QueueManager, Call, Timeout) ->
Start = now_in_milliseconds(),
ExpiresAt = expires(Timeout, Start),
try gen_server:call(QueueManager, {available_worker, ExpiresAt}, Timeout) of
{'EXIT', _, noproc} ->
noproc;
{'EXIT', Worker, Exit} ->
exit({Exit, {gen_server, call, [Worker, Call, Timeout]}});
{ok, Worker} ->
TimeLeft = time_left(ExpiresAt),
{ok, TimeLeft, Worker}
catch
_:{noproc, {gen_server, call, _}} ->
noproc;
_:{timeout, {gen_server, call, _}} ->
timeout
end.
inc_pending_tasks() ->
inc(pending_tasks).
dec_pending_tasks() ->
dec(pending_tasks).
inc(Key) ->
put(Key, get(Key) + 1).
dec(Key) ->
put(Key, get(Key) - 1).
-spec expires(timeout(), integer()) -> timeout().
expires(infinity, _) ->
infinity;
expires(Timeout, NowMs) ->
NowMs + Timeout.
-spec time_left(timeout()) -> timeout().
time_left(infinity) ->
infinity;
time_left(ExpiresAt) ->
ExpiresAt - now_in_milliseconds().
-spec is_expired(integer()) -> boolean().
is_expired(ExpiresAt) ->
ExpiresAt > now_in_milliseconds().
-spec now_in_milliseconds() -> integer().
now_in_milliseconds() ->
erlang:system_time(millisecond).
monitor_worker(Worker, Client, #state{monitors = Mons} = State) ->
Ref = monitor(process, Worker),
State#state{monitors = maps:put(Worker, {Ref, Client}, Mons)}.
queue_out(Clients, fifo) ->
queue:out(Clients);
queue_out(Clients, lifo) ->
queue:out_r(Clients).