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([initialised/1, selecting/2, new/1, new_with_initialiser/2, on_message/2, size/2, checkout_strategy/2, pid/1, apply/3, send/3, call/4, broadcast/2, shutdown/1, start/2, supervised/3]).
-export_type([checkout_strategy/0, initialised/2, builder/2, apply_error/0, 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 initialised(FJQ, FJR) :: {initialised,
FJQ,
gleam@option:option(gleam@erlang@process:selector(FJR))}.
-opaque builder(FJS, FJT) :: {builder,
integer(),
checkout_strategy(),
fun((gleam@erlang@process:subject(FJT)) -> {ok, initialised(FJS, FJT)} |
{error, binary()}),
integer(),
fun((FJS, FJT) -> gleam@otp@actor:next(FJS, FJT))}.
-type apply_error() :: no_resources_available.
-opaque pool(FJU) :: {pool, gleam@erlang@process:name(pool_msg(FJU))}.
-type state(FJV) :: {state,
gleam@deque:deque(worker(FJV)),
checkout_strategy(),
gleam@dict:dict(gleam@erlang@process:pid_(), live_worker(FJV)),
gleam@erlang@process:selector(pool_msg(FJV))}.
-type live_worker(FJW) :: {live_worker,
worker(FJW),
gleam@erlang@process:pid_(),
gleam@erlang@process:monitor()}.
-type pool_msg(FJX) :: {register, gleam@erlang@process:subject(FJX)} |
{check_in, worker(FJX), gleam@erlang@process:pid_()} |
{check_out,
gleam@erlang@process:subject({ok, worker(FJX)} | {error, apply_error()}),
gleam@erlang@process:pid_()} |
{worker_down, gleam@erlang@process:down()} |
{caller_down, gleam@erlang@process:down()} |
{broadcast, FJX}.
-type worker(FJY) :: {worker,
gleam@erlang@process:subject(FJY),
gleam@erlang@process:monitor()}.
-file("src/lifeguard.gleam", 35).
?DOC(" Create a new [`Initialised`](#Initialised) value with the given state and no selector.\n").
-spec initialised(FKD) -> initialised(FKD, any()).
initialised(State) ->
{initialised, State, none}.
-file("src/lifeguard.gleam", 46).
?DOC(
" Provide a selector for your worker to receive messages with. If your worker receives\n"
" a message that isn't selected, the message will be discarded and a warning will be\n"
" logged.\n"
"\n"
" If you don't provide a selector, the worker will receive messages from its own subject,\n"
" equivalent to `process.select(process.new_selector(), process.new_subject())`. If you\n"
" provide a selector, the default selector will be overwritten.\n"
).
-spec selecting(initialised(FKH, FKI), gleam@erlang@process:selector(FKI)) -> initialised(FKH, FKI).
selecting(Initialised, Selector) ->
_record = Initialised,
{initialised, erlang:element(2, _record), {some, Selector}}.
-file("src/lifeguard.gleam", 89).
?DOC(
" Create a new [`Builder`](#Builder) for creating a pool of worker actors.\n"
"\n"
" This API mimics the Actor API from [`gleam/otp/actor`](https://hexdocs.pm/gleam_otp/gleam/otp/actor.html),\n"
" so it should be familiar to anyone already using OTP with Gleam.\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(initial_state)\n"
" |> lifeguard.on_message(fn(msg, state) { actor.continue(state) })\n"
" |> lifeguard.size(10)\n"
" |> lifeguard.start(1000)\n"
" }\n"
" ```\n"
"\n"
" ### Default values\n"
"\n"
" | Config | Default |\n"
" |--------|---------|\n"
" | `on_message` | `fn(state, _) { actor.continue(state) }` |\n"
" | `size` | 10 |\n"
" | `checkout_strategy` | `FIFO` |\n"
).
-spec new(FKO) -> builder(FKO, any()).
new(State) ->
{builder,
10,
f_i_f_o,
fun(_) -> {ok, initialised(State)} end,
1000,
fun(State@1, _) -> gleam@otp@actor:continue(State@1) end}.
-file("src/lifeguard.gleam", 133).
?DOC(
" Create a new [`Builder`](#Builder) with a custom initialiser that runs before the\n"
" worker's message loop starts.\n"
"\n"
" The first argument is the number of milliseconds the initialiser is expected to\n"
" return within. The actor will be terminated if it does not complete within the\n"
" specified time, and the creation of the pool will fail.\n"
"\n"
" The initialiser is given the worker's default subject, which can optionally be\n"
" used to create a custom selector for the worker to receive messages. See the\n"
" [`selecting`](#selecting) function for more information.\n"
"\n"
" This API mimics the Actor API from [`gleam/otp/actor`](https://hexdocs.pm/gleam_otp/gleam/otp/actor.html),\n"
" so it should be familiar to anyone already using OTP with Gleam.\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(initial_state)\n"
" |> lifeguard.on_message(fn(msg, state) { actor.continue(state) })\n"
" |> lifeguard.size(10)\n"
" |> lifeguard.start(1000)\n"
" }\n"
" ```\n"
"\n"
" ### Default values\n"
"\n"
" | Config | Default |\n"
" |--------|---------|\n"
" | `on_message` | `fn(state, _) { actor.continue(state) }` |\n"
" | `size` | 10 |\n"
" | `checkout_strategy` | `FIFO` |\n"
).
-spec new_with_initialiser(
integer(),
fun((gleam@erlang@process:subject(FKS)) -> {ok, initialised(FKU, FKS)} |
{error, binary()})
) -> builder(FKU, FKS).
new_with_initialiser(Timeout, Initialiser) ->
{builder,
10,
f_i_f_o,
Initialiser,
Timeout,
fun(State, _) -> gleam@otp@actor:continue(State) end}.
-file("src/lifeguard.gleam", 149).
?DOC(
" Set the message handler for actors in the pool. This operates exactly like\n"
" [`gleam/otp/actor.on_message`](https://gleam.run/api/gleam/otp/actor/index.html#on_message).\n"
).
-spec on_message(
builder(FLB, FLC),
fun((FLB, FLC) -> gleam@otp@actor:next(FLB, FLC))
) -> builder(FLB, FLC).
on_message(Builder, Handler) ->
_record = Builder,
{builder,
erlang:element(2, _record),
erlang:element(3, _record),
erlang:element(4, _record),
erlang:element(5, _record),
Handler}.
-file("src/lifeguard.gleam", 157).
?DOC(" Set the number of actors in the pool. Defaults to 10.\n").
-spec size(builder(FLI, FLJ), integer()) -> builder(FLI, FLJ).
size(Builder, Size) ->
_record = Builder,
{builder,
Size,
erlang:element(3, _record),
erlang:element(4, _record),
erlang:element(5, _record),
erlang:element(6, _record)}.
-file("src/lifeguard.gleam", 165).
?DOC(" Set the order in which actors are checked out from the pool. Defaults to `FIFO`.\n").
-spec checkout_strategy(builder(FLO, FLP), checkout_strategy()) -> builder(FLO, FLP).
checkout_strategy(Builder, Checkout_strategy) ->
_record = Builder,
{builder,
erlang:element(2, _record),
Checkout_strategy,
erlang:element(4, _record),
erlang:element(5, _record),
erlang:element(6, _record)}.
-file("src/lifeguard.gleam", 283).
?DOC(
" Get the supervisor PID for a running pool.\n"
"\n"
" Returns an error if the pool is not running.\n"
).
-spec pid(pool(any())) -> {ok, gleam@erlang@process:pid_()} | {error, nil}.
pid(Pool) ->
gleam_erlang_ffi:process_named(erlang:element(2, Pool)).
-file("src/lifeguard.gleam", 287).
-spec check_out(pool(FMV), gleam@erlang@process:pid_(), integer()) -> {ok,
worker(FMV)} |
{error, apply_error()}.
check_out(Pool, Caller, Timeout) ->
gleam@erlang@process:call(
gleam@erlang@process:named_subject(erlang:element(2, Pool)),
Timeout,
fun(_capture) -> {check_out, _capture, Caller} end
).
-file("src/lifeguard.gleam", 295).
-spec check_in(pool(FNA), worker(FNA), gleam@erlang@process:pid_()) -> nil.
check_in(Pool, Worker, Caller) ->
gleam@erlang@process:send(
gleam@erlang@process:named_subject(erlang:element(2, Pool)),
{check_in, Worker, Caller}
).
-file("src/lifeguard.gleam", 300).
?DOC(false).
-spec apply(
pool(FNE),
integer(),
fun((gleam@erlang@process:subject(FNE)) -> FNH)
) -> {ok, FNH} | {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", 321).
?DOC(
" Send a message to a pooled worker. Equivalent to `process.send` using a pooled actor.\n"
"\n"
" ## Panics\n"
"\n"
" Like [`gleam/erlang/process.send`](https://hexdocs.pm/gleam_erlang/gleam/erlang/process.html#send),\n"
" this will panic if the pool is not running.\n"
).
-spec send(pool(FNK), FNK, integer()) -> {ok, nil} | {error, apply_error()}.
send(Pool, Msg, Checkout_timeout) ->
apply(
Pool,
Checkout_timeout,
fun(_capture) -> gleam@erlang@process:send(_capture, Msg) end
).
-file("src/lifeguard.gleam", 337).
?DOC(
" Send a message to a pooled actor and wait for a response. Equivalent to `process.call`\n"
" using a pooled worker.\n"
"\n"
" ## Panics\n"
"\n"
" Like [`gleam/erlang/process.call`](https://hexdocs.pm/gleam_erlang/gleam/erlang/process.html#call),\n"
" this will panic if the pool is not running, or if the worker crashes while\n"
" handling the message.\n"
).
-spec call(
pool(FNO),
fun((gleam@erlang@process:subject(FNQ)) -> FNO),
integer(),
integer()
) -> {ok, FNQ} | {error, apply_error()}.
call(Pool, Msg, Checkout_timeout, Call_timeout) ->
apply(
Pool,
Checkout_timeout,
fun(_capture) ->
gleam@erlang@process:call(_capture, Call_timeout, Msg)
end
).
-file("src/lifeguard.gleam", 352).
?DOC(
" Send a message to all pooled actors, regardless of checkout status.\n"
"\n"
" ## Panics\n"
"\n"
" Like [`gleam/erlang/process.send`](https://hexdocs.pm/gleam_erlang/gleam/erlang/process.html#send),\n"
" this will panic if the pool is not running.\n"
).
-spec broadcast(pool(FNU), FNU) -> nil.
broadcast(Pool, Msg) ->
_pipe = gleam@erlang@process:named_subject(erlang:element(2, Pool)),
gleam@erlang@process:send(_pipe, {broadcast, Msg}).
-file("src/lifeguard.gleam", 358).
?DOC(" Shut down a pool and all its workers. Fails if the pool is not currently running.\n").
-spec shutdown(pool(any())) -> {ok, nil} | {error, nil}.
shutdown(Pool) ->
_pipe = gleam_erlang_ffi:process_named(erlang:element(2, Pool)),
gleam@result:map(_pipe, fun gleam@erlang@process:send_exit/1).
-file("src/lifeguard.gleam", 395).
-spec handle_pool_message(state(FNZ), pool_msg(FNZ)) -> gleam@otp@actor:next(state(FNZ), pool_msg(FNZ)).
handle_pool_message(State, Msg) ->
case Msg of
{register, Worker_subject} ->
Pid_result = begin
_pipe = Worker_subject,
_pipe@1 = gleam@erlang@process:subject_owner(_pipe),
gleam@result:replace_error(
_pipe@1,
gleam@otp@actor:continue(State)
)
end,
case Pid_result of
{error, _} ->
gleam@otp@actor:continue(State);
{ok, Worker_pid} ->
Monitor = gleam@erlang@process:monitor(Worker_pid),
Selector = begin
_pipe@2 = erlang:element(5, State),
gleam@erlang@process:select_specific_monitor(
_pipe@2,
Monitor,
fun(Field@0) -> {worker_down, Field@0} end
)
end,
New_worker = {worker, Worker_subject, Monitor},
_pipe@3 = begin
_record = State,
{state,
gleam@deque:push_back(
erlang:element(2, State),
New_worker
),
erlang:element(3, _record),
erlang:element(4, _record),
Selector}
end,
_pipe@4 = gleam@otp@actor:continue(_pipe@3),
gleam@otp@actor:with_selector(_pipe@4, Selector)
end;
{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@5 = erlang:element(5, State),
gleam@erlang@process:deselect_specific_monitor(
_pipe@5,
erlang:element(4, Live_worker)
);
{error, _} ->
erlang:element(5, State)
end,
New_workers = gleam@deque:push_back(
erlang:element(2, State),
Worker
),
_pipe@6 = begin
_record@1 = State,
{state,
New_workers,
erlang:element(3, _record@1),
Live_workers,
Selector@1}
end,
_pipe@7 = gleam@otp@actor:continue(_pipe@6),
gleam@otp@actor:with_selector(_pipe@7, 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(Caller@1),
Selector@2 = begin
_pipe@8 = erlang:element(5, State),
gleam@erlang@process:select_specific_monitor(
_pipe@8,
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}),
_pipe@9 = begin
_record@2 = State,
{state,
New_workers@1,
erlang:element(3, _record@2),
Live_workers@1,
Selector@2}
end,
_pipe@10 = gleam@otp@actor:continue(_pipe@9),
gleam@otp@actor:with_selector(_pipe@10, Selector@2);
{error, _} ->
gleam@otp@actor:send(
Reply_to,
{error, no_resources_available}
),
gleam@otp@actor:continue(State)
end;
{caller_down, Process_down} ->
Caller_pid@1 = case Process_down of
{process_down, _, Caller_pid, _} -> Caller_pid;
_assert_fail ->
erlang:error(#{gleam_error => let_assert,
message => <<"Pattern match failed, no pattern matched the value."/utf8>>,
file => <<?FILEPATH/utf8>>,
module => <<"lifeguard"/utf8>>,
function => <<"handle_pool_message"/utf8>>,
line => 489,
value => _assert_fail,
start => 15585,
'end' => 15651,
pattern_start => 15596,
pattern_end => 15636})
end,
{Selector@4, Workers, Live_workers@3} = case gleam_stdlib:map_get(
erlang:element(4, State),
Caller_pid@1
) of
{ok, Live_worker@1} ->
Live_workers@2 = gleam@dict:delete(
erlang:element(4, State),
Caller_pid@1
),
gleam@erlang@process:demonitor_process(
erlang:element(4, Live_worker@1)
),
Selector@3 = begin
_pipe@11 = erlang:element(5, State),
gleam@erlang@process:deselect_specific_monitor(
_pipe@11,
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} ->
Worker_pid@2 = case Process_down@1 of
{process_down, _, Worker_pid@1, _} -> Worker_pid@1;
_assert_fail@1 ->
erlang:error(#{gleam_error => let_assert,
message => <<"Pattern match failed, no pattern matched the value."/utf8>>,
file => <<?FILEPATH/utf8>>,
module => <<"lifeguard"/utf8>>,
function => <<"handle_pool_message"/utf8>>,
line => 514,
value => _assert_fail@1,
start => 16501,
'end' => 16567,
pattern_start => 16512,
pattern_end => 16552})
end,
{Maybe_downed_worker, New_workers@3} = begin
_pipe@12 = erlang:element(2, State),
_pipe@13 = gleam@deque:to_list(_pipe@12),
gleam@list:partition(
_pipe@13,
fun(Worker@2) ->
gleam@erlang@process:subject_owner(
erlang:element(2, Worker@2)
)
=:= {ok, Worker_pid@2}
end
)
end,
{Downed_worker, Live_workers@4} = case Maybe_downed_worker of
[Worker@3] ->
{{some, Worker@3}, erlang:element(4, State)};
_ ->
case begin
_pipe@14 = maps:values(erlang:element(4, State)),
gleam@list:find(
_pipe@14,
fun(Lw) ->
gleam@erlang@process:subject_owner(
erlang:element(2, erlang:element(2, Lw))
)
=:= {ok, Worker_pid@2}
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@15 = erlang:element(5, State),
gleam@erlang@process:deselect_specific_monitor(
_pipe@15,
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@16 = New_workers@3,
gleam@deque:from_list(_pipe@16)
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@17 = maps:values(erlang:element(4, State)),
gleam@list:map(
_pipe@17,
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", 588).
-spec pool_spec(
builder(any(), FOE),
gleam@erlang@process:name(pool_msg(FOE)),
integer()
) -> gleam@otp@supervision:child_specification(gleam@erlang@process:subject(pool_msg(FOE))).
pool_spec(Builder, Pool_name, Init_timeout) ->
Pool_builder = begin
_pipe@4 = gleam@otp@actor:new_with_initialiser(
Init_timeout,
fun(Self) ->
Selector = begin
_pipe = gleam_erlang_ffi:new_selector(),
gleam@erlang@process:select(_pipe, Self)
end,
State = {state,
gleam@deque:new(),
erlang:element(3, Builder),
maps:new(),
Selector},
_pipe@1 = gleam@otp@actor:initialised(State),
_pipe@2 = gleam@otp@actor:selecting(_pipe@1, Selector),
_pipe@3 = gleam@otp@actor:returning(_pipe@2, Self),
{ok, _pipe@3}
end
),
_pipe@5 = gleam@otp@actor:on_message(_pipe@4, fun handle_pool_message/2),
gleam@otp@actor:named(_pipe@5, Pool_name)
end,
_pipe@6 = gleam@otp@supervision:worker(
fun() -> gleam@otp@actor:start(Pool_builder) end
),
gleam@otp@supervision:restart(_pipe@6, transient).
-file("src/lifeguard.gleam", 625).
-spec worker_spec(gleam@erlang@process:name(pool_msg(FOM)), builder(any(), FOM)) -> gleam@otp@supervision:child_specification(gleam@erlang@process:subject(FOM)).
worker_spec(Pool_name, Builder) ->
Worker_builder = begin
_pipe@4 = gleam@otp@actor:new_with_initialiser(
erlang:element(5, Builder),
fun(Self) ->
gleam@erlang@process:send(
gleam@erlang@process:named_subject(Pool_name),
{register, Self}
),
gleam@result:'try'(
(erlang:element(4, Builder))(Self),
fun(Init_data) ->
Selector = gleam@option:unwrap(
erlang:element(3, Init_data),
begin
_pipe = gleam_erlang_ffi:new_selector(),
gleam@erlang@process:select(_pipe, Self)
end
),
_pipe@1 = gleam@otp@actor:initialised(
erlang:element(2, Init_data)
),
_pipe@2 = gleam@otp@actor:selecting(_pipe@1, Selector),
_pipe@3 = gleam@otp@actor:returning(_pipe@2, Self),
{ok, _pipe@3}
end
)
end
),
gleam@otp@actor:on_message(_pipe@4, erlang:element(6, Builder))
end,
_pipe@5 = gleam@otp@supervision:worker(
fun() -> gleam@otp@actor:start(Worker_builder) end
),
gleam@otp@supervision:restart(_pipe@5, transient).
-file("src/lifeguard.gleam", 186).
?DOC(
" Start a pool supervision tree using the given [`Builder`](#Builder) 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_tree(
gleam@erlang@process:name(pool_msg(FLU)),
builder(any(), FLU),
integer()
) -> {ok, gleam@otp@actor:started(gleam@otp@static_supervisor:supervisor())} |
{error, gleam@otp@actor:start_error()}.
start_tree(Pool_name, Builder, Init_timeout) ->
Main_supervisor = gleam@otp@static_supervisor:new(rest_for_one),
Worker_supervisor = gleam@otp@static_supervisor:new(one_for_one),
Worker_supervisor_spec = gleam@otp@supervision:supervisor(
fun() ->
_pipe = gleam@list:repeat(<<""/utf8>>, erlang:element(2, Builder)),
_pipe@1 = gleam@list:fold(
_pipe,
Worker_supervisor,
fun(Worker_supervisor@1, _) ->
gleam@otp@static_supervisor:add(
Worker_supervisor@1,
worker_spec(Pool_name, Builder)
)
end
),
gleam@otp@static_supervisor:start(_pipe@1)
end
),
_pipe@2 = Main_supervisor,
_pipe@3 = gleam@otp@static_supervisor:add(
_pipe@2,
pool_spec(Builder, Pool_name, Init_timeout)
),
_pipe@4 = gleam@otp@static_supervisor:add(_pipe@3, Worker_supervisor_spec),
gleam@otp@static_supervisor:start(_pipe@4).
-file("src/lifeguard.gleam", 228).
?DOC(
" Start an unsupervised pool using the given [`Builder`](#Builder) and return a\n"
" [`Pool`](#Pool). In most cases, you should use the [`supervised`](#supervised)\n"
" function instead.\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(builder(any(), FME), integer()) -> {ok, pool(FME)} |
{error, gleam@otp@actor:start_error()}.
start(Builder, Init_timeout) ->
Pool_name = gleam_erlang_ffi:new_name(<<"lifeguard_pool"/utf8>>),
_pipe = start_tree(Pool_name, Builder, Init_timeout),
gleam@result:replace(_pipe, {pool, Pool_name}).
-file("src/lifeguard.gleam", 266).
?DOC(
" Return the [`ChildSpecification`](https://hexdocs.pm/gleam_otp/gleam/otp/supervision.html#ChildSpecification)\n"
" for creating a supervised worker pool.\n"
"\n"
" You must provide a selector to receive the [`Pool`](#Pool) value representing the\n"
" pool once it has started.\n"
"\n"
" ## Example\n"
"\n"
" ```gleam\n"
" import gleam/erlang/process\n"
" import gleam/otp/static_supervisor as supervisor\n"
" import lifeguard\n"
"\n"
" let pool_receiver = process.new_subject()\n"
"\n"
" let assert Ok(_started) =\n"
" supervisor.new(supervisor.OneForOne)\n"
" |> supervisor.add(\n"
" lifeguard.new(state)\n"
" |> lifeguard.supervised(pool_receiver, 1000)\n"
" )\n"
" |> supervisor.start\n"
"\n"
" let assert Ok(pool) =\n"
" process.receive(pool_receiver)\n"
"\n"
" let assert Ok(_) = lifeguard.send(pool, Message)\n"
" ```\n"
).
-spec supervised(
builder(any(), FML),
gleam@erlang@process:subject(pool(FML)),
integer()
) -> gleam@otp@supervision:child_specification(gleam@otp@static_supervisor:supervisor()).
supervised(Builder, Pool_subject, Init_timeout) ->
Pool_name = gleam_erlang_ffi:new_name(<<"lifeguard_pool"/utf8>>),
gleam@otp@supervision:supervisor(
fun() ->
gleam@result:'try'(
start_tree(Pool_name, Builder, Init_timeout),
fun(Supervisor) ->
gleam@erlang@process:send(Pool_subject, {pool, Pool_name}),
{ok, Supervisor}
end
)
end
).