Current section

Files

Jump to
lifeguard src lifeguard.erl
Raw

src/lifeguard.erl

-module(lifeguard).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch]).
-export([new/1, with_size/2, with_checkout_strategy/2, supervisor/1, apply/3, send/3, call/4, 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]).
-type checkout_strategy() :: fifo | lifo.
-opaque pool_config(HRG, HRH) :: {pool_config,
integer(),
spec(HRG, HRH),
checkout_strategy()}.
-type start_error() :: {pool_actor_start_error, gleam@otp@actor:start_error()} |
{worker_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(HRI, HRJ) :: {spec,
fun((gleam@erlang@process:selector(HRJ)) -> gleam@otp@actor:init_result(HRI, HRJ)),
integer(),
fun((HRJ, HRI) -> gleam@otp@actor:next(HRJ, HRI))}.
-opaque pool(HRK) :: {pool,
gleam@erlang@process:subject(pool_msg(HRK)),
gleam@erlang@process:pid_()}.
-type state(HRL) :: {state,
gleam@deque:deque(worker(HRL)),
checkout_strategy(),
gleam@dict:dict(gleam@erlang@process:pid_(), live_worker(HRL)),
gleam@erlang@process:selector(pool_msg(HRL))}.
-type live_worker(HRM) :: {live_worker,
worker(HRM),
gleam@erlang@process:pid_(),
gleam@erlang@process:process_monitor()}.
-type pool_msg(HRN) :: {register, gleam@erlang@process:subject(HRN)} |
{check_in, worker(HRN), gleam@erlang@process:pid_()} |
{check_out,
gleam@erlang@process:subject({ok, worker(HRN)} | {error, apply_error()}),
gleam@erlang@process:pid_()} |
{worker_down, gleam@erlang@process:process_down()} |
{caller_down, gleam@erlang@process:process_down()}.
-type worker(HRO) :: {worker,
gleam@erlang@process:subject(HRO),
gleam@erlang@process:process_monitor()}.
-file("/Users/isaac/repos/lifeguard/src/lifeguard.gleam", 57).
-spec new(spec(HRT, HRU)) -> pool_config(HRT, HRU).
new(Spec) ->
{pool_config, 10, Spec, fifo}.
-file("/Users/isaac/repos/lifeguard/src/lifeguard.gleam", 62).
-spec with_size(pool_config(HRZ, HSA), integer()) -> pool_config(HRZ, HSA).
with_size(Pool_config, Size) ->
erlang:setelement(2, Pool_config, Size).
-file("/Users/isaac/repos/lifeguard/src/lifeguard.gleam", 70).
-spec with_checkout_strategy(pool_config(HSF, HSG), checkout_strategy()) -> pool_config(HSF, HSG).
with_checkout_strategy(Pool_config, Checkout_strategy) ->
erlang:setelement(4, Pool_config, Checkout_strategy).
-file("/Users/isaac/repos/lifeguard/src/lifeguard.gleam", 198).
-spec supervisor(pool(any())) -> gleam@erlang@process:pid_().
supervisor(Pool) ->
erlang:element(3, Pool).
-file("/Users/isaac/repos/lifeguard/src/lifeguard.gleam", 202).
-spec check_out(pool(HSU), gleam@erlang@process:pid_(), integer()) -> {ok,
worker(HSU)} |
{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("/Users/isaac/repos/lifeguard/src/lifeguard.gleam", 212).
-spec check_in(pool(HSZ), worker(HSZ), gleam@erlang@process:pid_()) -> nil.
check_in(Pool, Worker, Caller) ->
gleam@erlang@process:send(
erlang:element(2, Pool),
{check_in, Worker, Caller}
).
-file("/Users/isaac/repos/lifeguard/src/lifeguard.gleam", 217).
-spec apply(
pool(HTD),
integer(),
fun((gleam@erlang@process:subject(HTD)) -> HTG)
) -> {ok, HTG} | {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("/Users/isaac/repos/lifeguard/src/lifeguard.gleam", 233).
-spec send(pool(HTJ), HTJ, 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("/Users/isaac/repos/lifeguard/src/lifeguard.gleam", 245).
-spec call(
pool(HTN),
fun((gleam@erlang@process:subject(HTP)) -> HTN),
integer(),
integer()
) -> {ok, HTP} | {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("/Users/isaac/repos/lifeguard/src/lifeguard.gleam", 264).
-spec shutdown(pool(any())) -> nil.
shutdown(Pool) ->
gleam@erlang@process:send_exit(erlang:element(3, Pool)).
-file("/Users/isaac/repos/lifeguard/src/lifeguard.gleam", 303).
-spec handle_pool_message(pool_msg(HTW), state(HTW)) -> gleam@otp@actor:next(pool_msg(HTW), state(HTW)).
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(
erlang:setelement(
5,
erlang:setelement(
2,
State,
gleam@deque:push_back(
erlang:element(2, State),
New_worker
)
),
Selector
)
),
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_ffi:demonitor(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(
erlang:setelement(
5,
erlang:setelement(
4,
erlang:setelement(2, State, New_workers),
Live_workers
),
Selector@1
)
),
Selector@1
);
{check_out, Reply_to, Caller@1} ->
Get_result = case erlang:element(3, State) of
fifo ->
gleam@deque:pop_front(erlang:element(2, State));
lifo ->
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(
erlang:setelement(
5,
erlang:setelement(
4,
erlang:setelement(2, State, New_workers@1),
Live_workers@1
),
Selector@2
)
),
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_ffi:demonitor(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(
erlang:setelement(
2,
erlang:setelement(
4,
erlang:setelement(5, State, Selector@4),
Live_workers@3
),
Workers
)
),
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_ffi:demonitor(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(
erlang:setelement(
2,
erlang:setelement(
5,
erlang:setelement(4, State, Live_workers@4),
Selector@5
),
begin
_pipe@9 = New_workers@3,
gleam@deque:from_list(_pipe@9)
end
)
),
Selector@5
)
end.
-file("/Users/isaac/repos/lifeguard/src/lifeguard.gleam", 471).
-spec pool_spec(pool_config(any(), HUB), integer()) -> gleam@otp@actor:spec(state(HUB), pool_msg(HUB)).
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("/Users/isaac/repos/lifeguard/src/lifeguard.gleam", 500).
-spec worker_spec(gleam@erlang@process:subject(pool_msg(HUI)), spec(HUL, HUI)) -> gleam@otp@actor:spec(HUL, HUI).
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("/Users/isaac/repos/lifeguard/src/lifeguard.gleam", 124).
-spec start(pool_config(any(), HSM), integer()) -> {ok, pool(HSM)} |
{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) ->
Workers_result = begin
_pipe@1 = gleam@list:repeat(
<<""/utf8>>,
erlang:element(2, Pool_config)
),
gleam@list:try_map(
_pipe@1,
fun(_) ->
gleam@result:'try'(
begin
_pipe@2 = gleam@otp@actor:start_spec(
worker_spec(
Pool_subject,
erlang:element(3, Pool_config)
)
),
gleam@result:map_error(
_pipe@2,
fun(Field@0) -> {worker_start_error, Field@0} end
)
end,
fun(Subject) -> {ok, Subject} end
)
end
)
end,
gleam@result:'try'(
Workers_result,
fun(Workers) ->
Worker_supervisor_result = begin
_pipe@3 = Workers,
_pipe@6 = gleam@list:index_fold(
_pipe@3,
Worker_supervisor,
fun(Worker_supervisor@1, Actor, Idx) ->
gleam@otp@static_supervisor:add(
Worker_supervisor@1,
begin
_pipe@5 = gleam@otp@static_supervisor:worker_child(
<<"worker_"/utf8,
(erlang:integer_to_binary(Idx))/binary>>,
fun() ->
_pipe@4 = gleam@erlang@process:subject_owner(
Actor
),
{ok, _pipe@4}
end
),
gleam@otp@static_supervisor:restart(
_pipe@5,
transient
)
end
)
end
),
_pipe@7 = gleam@otp@static_supervisor:start_link(
_pipe@6
),
gleam@result:map_error(
_pipe@7,
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@10 = gleam@otp@static_supervisor:add(
Main_supervisor,
begin
_pipe@9 = gleam@otp@static_supervisor:worker_child(
<<"pool"/utf8>>,
fun() ->
_pipe@8 = gleam@erlang@process:subject_owner(
Pool_subject
),
{ok, _pipe@8}
end
),
gleam@otp@static_supervisor:restart(
_pipe@9,
transient
)
end
),
_pipe@12 = gleam@otp@static_supervisor:add(
_pipe@10,
begin
_pipe@11 = gleam@otp@static_supervisor:supervisor_child(
<<"worker_supervisor"/utf8>>,
fun() ->
{ok, Worker_supervisor@2}
end
),
gleam@otp@static_supervisor:restart(
_pipe@11,
transient
)
end
),
_pipe@13 = gleam@otp@static_supervisor:start_link(
_pipe@12
),
gleam@result:map_error(
_pipe@13,
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
)
end
).