Current section
Files
Jump to
Current section
Files
src/lifeguard.erl
-module(lifeguard).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch]).
-define(FILEPATH, "src/lifeguard.gleam").
-export([new/1, with_size/2, with_checkout_strategy/2, supervisor/1, apply/3, send/3, call/4, broadcast/2, shutdown/1, start/2]).
-export_type([checkout_strategy/0, pool_config/2, start_error/0, apply_error/0, spec/2, pool/1, state/1, live_worker/1, pool_msg/1, worker/1]).
-if(?OTP_RELEASE >= 27).
-define(MODULEDOC(Str), -moduledoc(Str)).
-define(DOC(Str), -doc(Str)).
-else.
-define(MODULEDOC(Str), -compile([])).
-define(DOC(Str), -compile([])).
-endif.
-type checkout_strategy() :: f_i_f_o | l_i_f_o.
-opaque pool_config(GGT, GGU) :: {pool_config,
integer(),
spec(GGT, GGU),
checkout_strategy()}.
-type start_error() :: {pool_actor_start_error, gleam@otp@actor:start_error()} |
{pool_supervisor_start_error, gleam@dynamic:dynamic_()} |
{worker_supervisor_start_error, gleam@dynamic:dynamic_()}.
-type apply_error() :: no_resources_available |
check_out_timeout |
worker_call_timeout |
{worker_crashed, gleam@dynamic:dynamic_()}.
-type spec(GGV, GGW) :: {spec,
fun((gleam@erlang@process:selector(GGW)) -> gleam@otp@actor:init_result(GGV, GGW)),
integer(),
fun((GGW, GGV) -> gleam@otp@actor:next(GGW, GGV))}.
-opaque pool(GGX) :: {pool,
gleam@erlang@process:subject(pool_msg(GGX)),
gleam@erlang@process:pid_()}.
-type state(GGY) :: {state,
gleam@deque:deque(worker(GGY)),
checkout_strategy(),
gleam@dict:dict(gleam@erlang@process:pid_(), live_worker(GGY)),
gleam@erlang@process:selector(pool_msg(GGY))}.
-type live_worker(GGZ) :: {live_worker,
worker(GGZ),
gleam@erlang@process:pid_(),
gleam@erlang@process:process_monitor()}.
-type pool_msg(GHA) :: {register, gleam@erlang@process:subject(GHA)} |
{check_in, worker(GHA), gleam@erlang@process:pid_()} |
{check_out,
gleam@erlang@process:subject({ok, worker(GHA)} | {error, apply_error()}),
gleam@erlang@process:pid_()} |
{worker_down, gleam@erlang@process:process_down()} |
{caller_down, gleam@erlang@process:process_down()} |
{broadcast, GHA}.
-type worker(GHB) :: {worker,
gleam@erlang@process:subject(GHB),
gleam@erlang@process:process_monitor()}.
-file("src/lifeguard.gleam", 57).
?DOC(
" Create a new [`PoolConfig`](#PoolConfig) for creating a pool of actors.\n"
"\n"
" ```gleam\n"
" import lifeguard\n"
"\n"
" pub fn main() {\n"
" // Create a pool of 10 actors that do nothing.\n"
" let assert Ok(pool) =\n"
" lifeguard.new(\n"
" lifeguard.Spec(\n"
" init_timeout: 1000,\n"
" init: fn(selector) { actor.Ready(state: Nil, selector:) },\n"
" loop: fn(msg, state) { actor.continue(state) },\n"
" )\n"
" )\n"
" |> lifeguard.with_size(10)\n"
" |> lifeguard.start(1000)\n"
" }\n"
" ```\n"
"\n"
" ### Default values\n"
"\n"
" | Config | Default |\n"
" |--------|---------|\n"
" | `size` | 10 |\n"
" | `checkout_strategy` | `FIFO` |\n"
).
-spec new(spec(GHG, GHH)) -> pool_config(GHG, GHH).
new(Spec) ->
{pool_config, 10, Spec, f_i_f_o}.
-file("src/lifeguard.gleam", 62).
?DOC(" Set the number of actors in the pool. Defaults to 10.\n").
-spec with_size(pool_config(GHM, GHN), integer()) -> pool_config(GHM, GHN).
with_size(Pool_config, Size) ->
_record = Pool_config,
{pool_config, Size, erlang:element(3, _record), erlang:element(4, _record)}.
-file("src/lifeguard.gleam", 70).
?DOC(" Set the order in which actors are checked out from the pool. Defaults to `FIFO`.\n").
-spec with_checkout_strategy(pool_config(GHS, GHT), checkout_strategy()) -> pool_config(GHS, GHT).
with_checkout_strategy(Pool_config, Checkout_strategy) ->
_record = Pool_config,
{pool_config,
erlang:element(2, _record),
erlang:element(3, _record),
Checkout_strategy}.
-file("src/lifeguard.gleam", 191).
?DOC(" Get the supervisor PID for a running pool.\n").
-spec supervisor(pool(any())) -> gleam@erlang@process:pid_().
supervisor(Pool) ->
erlang:element(3, Pool).
-file("src/lifeguard.gleam", 195).
-spec check_out(pool(GIH), gleam@erlang@process:pid_(), integer()) -> {ok,
worker(GIH)} |
{error, apply_error()}.
check_out(Pool, Caller, Timeout) ->
_pipe = gleam@erlang@process:try_call(
erlang:element(2, Pool),
fun(_capture) -> {check_out, _capture, Caller} end,
Timeout
),
_pipe@1 = gleam@result:replace_error(_pipe, check_out_timeout),
gleam@result:flatten(_pipe@1).
-file("src/lifeguard.gleam", 205).
-spec check_in(pool(GIM), worker(GIM), gleam@erlang@process:pid_()) -> nil.
check_in(Pool, Worker, Caller) ->
gleam@erlang@process:send(
erlang:element(2, Pool),
{check_in, Worker, Caller}
).
-file("src/lifeguard.gleam", 210).
?DOC(false).
-spec apply(
pool(GIQ),
integer(),
fun((gleam@erlang@process:subject(GIQ)) -> GIT)
) -> {ok, GIT} | {error, apply_error()}.
apply(Pool, Timeout, Next) ->
Self = erlang:self(),
gleam@result:'try'(
check_out(Pool, Self, Timeout),
fun(Worker) ->
Result = Next(erlang:element(2, Worker)),
check_in(Pool, Worker, Self),
{ok, Result}
end
).
-file("src/lifeguard.gleam", 226).
?DOC(" Send a message to a pooled actor. Equivalent to `process.send` using a pooled actor.\n").
-spec send(pool(GIW), GIW, integer()) -> {ok, nil} | {error, apply_error()}.
send(Pool, Msg, Checkout_timeout) ->
apply(
Pool,
Checkout_timeout,
fun(Subject) -> gleam@erlang@process:send(Subject, Msg) end
).
-file("src/lifeguard.gleam", 238).
?DOC(
" Send a message to a pooled actor and wait for a response. Equivalent to `process.call`\n"
" using a pooled actor.\n"
).
-spec call(
pool(GJA),
fun((gleam@erlang@process:subject(GJC)) -> GJA),
integer(),
integer()
) -> {ok, GJC} | {error, apply_error()}.
call(Pool, Msg, Checkout_timeout, Call_timeout) ->
_pipe@1 = apply(
Pool,
Checkout_timeout,
fun(Subject) ->
_pipe = gleam@erlang@process:try_call(Subject, Msg, Call_timeout),
gleam@result:map_error(_pipe, fun(Err) -> case Err of
call_timeout ->
worker_call_timeout;
{callee_down, Reason} ->
{worker_crashed, Reason}
end end)
end
),
gleam@result:flatten(_pipe@1).
-file("src/lifeguard.gleam", 257).
?DOC(" Send a message to all pooled actors, regardless of checkout status.\n").
-spec broadcast(pool(GJG), GJG) -> nil.
broadcast(Pool, Msg) ->
gleam@erlang@process:send(erlang:element(2, Pool), {broadcast, Msg}).
-file("src/lifeguard.gleam", 262).
?DOC(" Shut down a pool and all its workers.\n").
-spec shutdown(pool(any())) -> nil.
shutdown(Pool) ->
gleam@erlang@process:send_exit(erlang:element(3, Pool)).
-file("src/lifeguard.gleam", 302).
-spec handle_pool_message(pool_msg(GJL), state(GJL)) -> gleam@otp@actor:next(pool_msg(GJL), state(GJL)).
handle_pool_message(Msg, State) ->
case Msg of
{register, Worker_subject} ->
Monitor = gleam@erlang@process:monitor_process(
begin
_pipe = Worker_subject,
gleam@erlang@process:subject_owner(_pipe)
end
),
Selector = begin
_pipe@1 = erlang:element(5, State),
gleam@erlang@process:selecting_process_down(
_pipe@1,
Monitor,
fun(Field@0) -> {worker_down, Field@0} end
)
end,
New_worker = {worker, Worker_subject, Monitor},
gleam@otp@actor:with_selector(
gleam@otp@actor:continue(
begin
_record = State,
{state,
gleam@deque:push_back(
erlang:element(2, State),
New_worker
),
erlang:element(3, _record),
erlang:element(4, _record),
Selector}
end
),
Selector
);
{check_in, Worker, Caller} ->
Caller_live_worker = gleam_stdlib:map_get(
erlang:element(4, State),
Caller
),
Live_workers = gleam@dict:delete(erlang:element(4, State), Caller),
Selector@1 = case Caller_live_worker of
{ok, Live_worker} ->
gleam@erlang@process:demonitor_process(
erlang:element(4, Live_worker)
),
_pipe@2 = erlang:element(5, State),
gleam@erlang@process:deselecting_process_down(
_pipe@2,
erlang:element(4, Live_worker)
);
{error, _} ->
erlang:element(5, State)
end,
New_workers = gleam@deque:push_back(
erlang:element(2, State),
Worker
),
gleam@otp@actor:with_selector(
gleam@otp@actor:continue(
begin
_record@1 = State,
{state,
New_workers,
erlang:element(3, _record@1),
Live_workers,
Selector@1}
end
),
Selector@1
);
{check_out, Reply_to, Caller@1} ->
Get_result = case erlang:element(3, State) of
f_i_f_o ->
gleam@deque:pop_front(erlang:element(2, State));
l_i_f_o ->
gleam@deque:pop_back(erlang:element(2, State))
end,
case Get_result of
{ok, {Worker@1, New_workers@1}} ->
Caller_monitor = gleam@erlang@process:monitor_process(
Caller@1
),
Selector@2 = begin
_pipe@3 = erlang:element(5, State),
gleam@erlang@process:selecting_process_down(
_pipe@3,
Caller_monitor,
fun(Field@0) -> {caller_down, Field@0} end
)
end,
Live_workers@1 = gleam@dict:insert(
erlang:element(4, State),
Caller@1,
{live_worker, Worker@1, Caller@1, Caller_monitor}
),
gleam@otp@actor:send(Reply_to, {ok, Worker@1}),
gleam@otp@actor:with_selector(
gleam@otp@actor:continue(
begin
_record@2 = State,
{state,
New_workers@1,
erlang:element(3, _record@2),
Live_workers@1,
Selector@2}
end
),
Selector@2
);
{error, _} ->
gleam@otp@actor:send(
Reply_to,
{error, no_resources_available}
),
gleam@otp@actor:continue(State)
end;
{caller_down, Process_down} ->
{Selector@4, Workers, Live_workers@3} = case gleam_stdlib:map_get(
erlang:element(4, State),
erlang:element(2, Process_down)
) of
{ok, Live_worker@1} ->
Live_workers@2 = gleam@dict:delete(
erlang:element(4, State),
erlang:element(2, Process_down)
),
gleam@erlang@process:demonitor_process(
erlang:element(4, Live_worker@1)
),
Selector@3 = begin
_pipe@4 = erlang:element(5, State),
gleam@erlang@process:deselecting_process_down(
_pipe@4,
erlang:element(4, Live_worker@1)
)
end,
New_workers@2 = gleam@deque:push_back(
erlang:element(2, State),
erlang:element(2, Live_worker@1)
),
{Selector@3, New_workers@2, Live_workers@2};
{error, _} ->
{erlang:element(5, State),
erlang:element(2, State),
erlang:element(4, State)}
end,
gleam@otp@actor:with_selector(
gleam@otp@actor:continue(
begin
_record@3 = State,
{state,
Workers,
erlang:element(3, _record@3),
Live_workers@3,
Selector@4}
end
),
Selector@4
);
{worker_down, Process_down@1} ->
{Maybe_downed_worker, New_workers@3} = begin
_pipe@5 = erlang:element(2, State),
_pipe@6 = gleam@deque:to_list(_pipe@5),
gleam@list:partition(
_pipe@6,
fun(Worker@2) ->
gleam@erlang@process:subject_owner(
erlang:element(2, Worker@2)
)
=:= erlang:element(2, Process_down@1)
end
)
end,
{Downed_worker, Live_workers@4} = case Maybe_downed_worker of
[Worker@3] ->
{{some, Worker@3}, erlang:element(4, State)};
_ ->
case begin
_pipe@7 = maps:values(erlang:element(4, State)),
gleam@list:find(
_pipe@7,
fun(Lw) ->
gleam@erlang@process:subject_owner(
erlang:element(2, erlang:element(2, Lw))
)
=:= erlang:element(2, Process_down@1)
end
)
end of
{ok, Live_worker@2} ->
{{some, erlang:element(2, Live_worker@2)},
gleam@dict:delete(
erlang:element(4, State),
erlang:element(3, Live_worker@2)
)};
{error, nil} ->
{none, erlang:element(4, State)}
end
end,
Selector@5 = case Downed_worker of
{some, Worker@4} ->
gleam@erlang@process:demonitor_process(
erlang:element(3, Worker@4)
),
_pipe@8 = erlang:element(5, State),
gleam@erlang@process:deselecting_process_down(
_pipe@8,
erlang:element(3, Worker@4)
);
_ ->
erlang:element(5, State)
end,
gleam@otp@actor:with_selector(
gleam@otp@actor:continue(
begin
_record@4 = State,
{state,
begin
_pipe@9 = New_workers@3,
gleam@deque:from_list(_pipe@9)
end,
erlang:element(3, _record@4),
Live_workers@4,
Selector@5}
end
),
Selector@5
);
{broadcast, Msg_to_send} ->
Workers@1 = lists:append(
gleam@deque:to_list(erlang:element(2, State)),
begin
_pipe@10 = maps:values(erlang:element(4, State)),
gleam@list:map(
_pipe@10,
fun(Live_worker@3) ->
erlang:element(2, Live_worker@3)
end
)
end
),
gleam@list:each(
Workers@1,
fun(Worker@5) ->
gleam@erlang@process:send(
erlang:element(2, Worker@5),
Msg_to_send
)
end
),
gleam@otp@actor:continue(State)
end.
-file("src/lifeguard.gleam", 484).
-spec pool_spec(pool_config(any(), GJQ), integer()) -> gleam@otp@actor:spec(state(GJQ), pool_msg(GJQ)).
pool_spec(Pool_config, Init_timeout) ->
{spec,
fun() ->
Self = gleam@erlang@process:new_subject(),
Selector = begin
_pipe = gleam_erlang_ffi:new_selector(),
gleam@erlang@process:selecting(
_pipe,
Self,
fun gleam@function:identity/1
)
end,
State = {state,
gleam@deque:new(),
erlang:element(4, Pool_config),
maps:new(),
Selector},
{ready, State, Selector}
end,
Init_timeout,
fun handle_pool_message/2}.
-file("src/lifeguard.gleam", 513).
-spec worker_spec(gleam@erlang@process:subject(pool_msg(GJX)), spec(GKA, GJX)) -> gleam@otp@actor:spec(GKA, GJX).
worker_spec(Pool_subject, Spec) ->
{spec,
fun() ->
Self = gleam@erlang@process:new_subject(),
gleam@erlang@process:send(Pool_subject, {register, Self}),
Selector = begin
_pipe = gleam_erlang_ffi:new_selector(),
gleam@erlang@process:selecting(
_pipe,
Self,
fun gleam@function:identity/1
)
end,
(erlang:element(2, Spec))(Selector)
end,
erlang:element(3, Spec),
erlang:element(4, Spec)}.
-file("src/lifeguard.gleam", 127).
?DOC(
" Start a pool supervision tree using the given [`PoolConfig`](#PoolConfig) and return a\n"
" [`Pool`](#Pool).\n"
"\n"
" Note: this function mimics the behaviour of `supervisor:start_link` and\n"
" `gleam/otp/static_supervisor`'s `start_link` function and will exit the process if\n"
" any of the workers fail to start.\n"
).
-spec start(pool_config(any(), GHZ), integer()) -> {ok, pool(GHZ)} |
{error, start_error()}.
start(Pool_config, Init_timeout) ->
Main_supervisor = gleam@otp@static_supervisor:new(rest_for_one),
Worker_supervisor = gleam@otp@static_supervisor:new(one_for_one),
Pool_start_result = begin
_pipe = gleam@otp@actor:start_spec(pool_spec(Pool_config, Init_timeout)),
gleam@result:map_error(
_pipe,
fun(Field@0) -> {pool_actor_start_error, Field@0} end
)
end,
gleam@result:'try'(
Pool_start_result,
fun(Pool_subject) ->
Worker_supervisor_result = begin
_pipe@1 = gleam@list:repeat(
<<""/utf8>>,
erlang:element(2, Pool_config)
),
_pipe@5 = gleam@list:index_fold(
_pipe@1,
Worker_supervisor,
fun(Worker_supervisor@1, _, Idx) ->
gleam@otp@static_supervisor:add(
Worker_supervisor@1,
begin
_pipe@4 = gleam@otp@static_supervisor:worker_child(
<<"worker_"/utf8,
(erlang:integer_to_binary(Idx))/binary>>,
fun() ->
_pipe@2 = worker_spec(
Pool_subject,
erlang:element(3, Pool_config)
),
_pipe@3 = gleam@otp@actor:start_spec(
_pipe@2
),
gleam@result:map(
_pipe@3,
fun gleam@erlang@process:subject_owner/1
)
end
),
gleam@otp@static_supervisor:restart(
_pipe@4,
transient
)
end
)
end
),
_pipe@6 = gleam@otp@static_supervisor:start_link(_pipe@5),
gleam@result:map_error(
_pipe@6,
fun(Field@0) -> {worker_supervisor_start_error, Field@0} end
)
end,
gleam@result:'try'(
Worker_supervisor_result,
fun(Worker_supervisor@2) ->
Main_supervisor_result = begin
_pipe@9 = gleam@otp@static_supervisor:add(
Main_supervisor,
begin
_pipe@8 = gleam@otp@static_supervisor:worker_child(
<<"pool"/utf8>>,
fun() ->
_pipe@7 = gleam@erlang@process:subject_owner(
Pool_subject
),
{ok, _pipe@7}
end
),
gleam@otp@static_supervisor:restart(
_pipe@8,
transient
)
end
),
_pipe@11 = gleam@otp@static_supervisor:add(
_pipe@9,
begin
_pipe@10 = gleam@otp@static_supervisor:supervisor_child(
<<"worker_supervisor"/utf8>>,
fun() -> {ok, Worker_supervisor@2} end
),
gleam@otp@static_supervisor:restart(
_pipe@10,
transient
)
end
),
_pipe@12 = gleam@otp@static_supervisor:start_link(
_pipe@11
),
gleam@result:map_error(
_pipe@12,
fun(Field@0) -> {pool_supervisor_start_error, Field@0} end
)
end,
gleam@result:'try'(
Main_supervisor_result,
fun(Main_supervisor@1) ->
{ok, {pool, Pool_subject, Main_supervisor@1}}
end
)
end
)
end
).