Current section
Files
Jump to
Current section
Files
src/bath.erl
-module(bath).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch, inline]).
-define(FILEPATH, "src/bath.gleam").
-export([new/1, size/2, name/2, on_shutdown/2, checkout_strategy/2, creation_strategy/2, log_errors/2, keep/0, discard/0, returning/2, apply/3, apply_blocking/3, shutdown/3, try_map_returning/2, supervised_map/3, supervised/2, start/2]).
-export_type([checkout_strategy/0, creation_strategy/0, builder/1, apply_error/0, shutdown_error/0, next/1, state/1, live_resource/1, waiting_request/1, msg/1, serve_result/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.
-type creation_strategy() :: lazy | eager.
-opaque builder(FDL) :: {builder,
gleam@option:option(gleam@erlang@process:name(msg(FDL))),
integer(),
fun(() -> {ok, FDL} | {error, binary()}),
fun((FDL) -> nil),
checkout_strategy(),
creation_strategy(),
boolean()}.
-type apply_error() :: no_resources_available |
{check_out_resource_create_error, binary()} |
pool_shutting_down.
-type shutdown_error() :: resources_in_use.
-opaque next(FDM) :: {keep, FDM} | {discard, FDM}.
-opaque state(FDN) :: {state,
checkout_strategy(),
creation_strategy(),
integer(),
fun(() -> {ok, FDN} | {error, binary()}),
fun((FDN) -> nil),
gleam@deque:deque(FDN),
integer(),
gleam@dict:dict(gleam@erlang@process:pid_(), live_resource(FDN)),
gleam@deque:deque(waiting_request(FDN)),
gleam@erlang@process:selector(msg(FDN)),
boolean()}.
-type live_resource(FDO) :: {live_resource, FDO, gleam@erlang@process:monitor()}.
-type waiting_request(FDP) :: {waiting_request,
gleam@erlang@process:subject({ok, FDP} | {error, apply_error()}),
gleam@erlang@process:pid_(),
gleam@erlang@process:monitor()}.
-opaque msg(FDQ) :: {check_in, FDQ, gleam@erlang@process:pid_(), next(nil)} |
{check_out,
gleam@erlang@process:subject({ok, FDQ} | {error, apply_error()}),
gleam@erlang@process:pid_()} |
{check_out_blocking,
gleam@erlang@process:subject({ok, FDQ} | {error, apply_error()}),
gleam@erlang@process:pid_()} |
{pool_exit, gleam@erlang@process:exit_message()} |
{caller_down, gleam@erlang@process:down()} |
{waiter_down, gleam@erlang@process:down()} |
{shutdown,
gleam@erlang@process:subject({ok, nil} | {error, shutdown_error()}),
boolean()}.
-type serve_result(FDR) :: {served,
gleam@dict:dict(gleam@erlang@process:pid_(), live_resource(FDR)),
gleam@deque:deque(waiting_request(FDR)),
gleam@erlang@process:selector(msg(FDR)),
integer()} |
{serve_failed,
gleam@deque:deque(waiting_request(FDR)),
gleam@erlang@process:selector(msg(FDR)),
integer()} |
none_waiting.
-file("src/bath.gleam", 68).
?DOC(
" Create a new [`Builder`](#Builder) for creating a pool of resources.\n"
"\n"
" ```gleam\n"
" import bath\n"
" import fake_db\n"
"\n"
" pub fn main() {\n"
" // Create a pool of 10 connections to some fictional database.\n"
" let assert Ok(pool) =\n"
" bath.new(fn() { fake_db.connect() })\n"
" |> bath.size(10)\n"
" |> bath.start(1000)\n"
" }\n"
" ```\n"
"\n"
" ### Default values\n"
"\n"
" | Config | Default |\n"
" |--------|---------|\n"
" | `name` | `option.None` |\n"
" | `size` | 10 |\n"
" | `shutdown_resource` | `fn(_resource) { Nil }` |\n"
" | `checkout_strategy` | `FIFO` |\n"
" | `creation_strategy` | `Lazy` |\n"
" | `log_errors` | `False` |\n"
).
-spec new(fun(() -> {ok, FDW} | {error, binary()})) -> builder(FDW).
new(Create_resource) ->
{builder,
none,
10,
Create_resource,
fun(_) -> nil end,
f_i_f_o,
lazy,
false}.
-file("src/bath.gleam", 83).
?DOC(" Set the pool size. Defaults to 10. Will be clamped to a minimum of 1.\n").
-spec size(builder(FEA), integer()) -> builder(FEA).
size(Builder, Size) ->
{builder,
erlang:element(2, Builder),
gleam@int:max(Size, 1),
erlang:element(4, Builder),
erlang:element(5, Builder),
erlang:element(6, Builder),
erlang:element(7, Builder),
erlang:element(8, Builder)}.
-file("src/bath.gleam", 94).
?DOC(
" Set the name for the pool process. Defaults to `None`.\n"
"\n"
" You will need to provide a name if you plan on using the pool under a static\n"
" supervisor.\n"
).
-spec name(builder(FED), gleam@erlang@process:name(msg(FED))) -> builder(FED).
name(Builder, Name) ->
{builder,
{some, Name},
erlang:element(3, Builder),
erlang:element(4, Builder),
erlang:element(5, Builder),
erlang:element(6, Builder),
erlang:element(7, Builder),
erlang:element(8, Builder)}.
-file("src/bath.gleam", 102).
?DOC(" Set a shutdown function to be run for each resource when the pool exits.\n").
-spec on_shutdown(builder(FEI), fun((FEI) -> nil)) -> builder(FEI).
on_shutdown(Builder, Shutdown_resource) ->
{builder,
erlang:element(2, Builder),
erlang:element(3, Builder),
erlang:element(4, Builder),
Shutdown_resource,
erlang:element(6, Builder),
erlang:element(7, Builder),
erlang:element(8, Builder)}.
-file("src/bath.gleam", 110).
?DOC(" Change the checkout strategy for the pool. Defaults to `FIFO`.\n").
-spec checkout_strategy(builder(FEL), checkout_strategy()) -> builder(FEL).
checkout_strategy(Builder, Checkout_strategy) ->
{builder,
erlang:element(2, Builder),
erlang:element(3, Builder),
erlang:element(4, Builder),
erlang:element(5, Builder),
Checkout_strategy,
erlang:element(7, Builder),
erlang:element(8, Builder)}.
-file("src/bath.gleam", 118).
?DOC(" Change the resource creation strategy for the pool. Defaults to `Lazy`.\n").
-spec creation_strategy(builder(FEO), creation_strategy()) -> builder(FEO).
creation_strategy(Builder, Creation_strategy) ->
{builder,
erlang:element(2, Builder),
erlang:element(3, Builder),
erlang:element(4, Builder),
erlang:element(5, Builder),
erlang:element(6, Builder),
Creation_strategy,
erlang:element(8, Builder)}.
-file("src/bath.gleam", 126).
?DOC(" Set whether the pool logs errors when resources fail to create.\n").
-spec log_errors(builder(FER), boolean()) -> builder(FER).
log_errors(Builder, Log_errors) ->
{builder,
erlang:element(2, Builder),
erlang:element(3, Builder),
erlang:element(4, Builder),
erlang:element(5, Builder),
erlang:element(6, Builder),
erlang:element(7, Builder),
Log_errors}.
-file("src/bath.gleam", 226).
?DOC(" Instruct Bath to keep the checked out resource, returning it to the pool.\n").
-spec keep() -> next(nil).
keep() ->
{keep, nil}.
-file("src/bath.gleam", 235).
?DOC(
" Instruct Bath to discard the checked out resource, running the shutdown function on\n"
" it.\n"
"\n"
" Discarded resources will be recreated lazily, regardless of the pool's creation\n"
" strategy.\n"
).
-spec discard() -> next(nil).
discard() ->
{discard, nil}.
-file("src/bath.gleam", 240).
?DOC(" Return a value from a use of [`apply`](#apply).\n").
-spec returning(next(any()), FFN) -> next(FFN).
returning(Next, Value) ->
case Next of
{keep, _} ->
{keep, Value};
{discard, _} ->
{discard, Value}
end.
-file("src/bath.gleam", 250).
?DOC(
" Checks out a resource from the pool, sending the caller Pid for the pool to\n"
" monitor in case the client dies. This allows the pool to create a new resource\n"
" later if required.\n"
).
-spec check_out(
gleam@erlang@process:subject(msg(FFP)),
gleam@erlang@process:pid_(),
integer()
) -> {ok, FFP} | {error, apply_error()}.
check_out(Pool, Caller, Timeout) ->
gleam@erlang@process:call(
Pool,
Timeout,
fun(_capture) -> {check_out, _capture, Caller} end
).
-file("src/bath.gleam", 258).
-spec check_in(
gleam@erlang@process:subject(msg(FFU)),
FFU,
gleam@erlang@process:pid_(),
next(nil)
) -> nil.
check_in(Pool, Resource, Caller, Next) ->
gleam@erlang@process:send(Pool, {check_in, Resource, Caller, Next}).
-file("src/bath.gleam", 283).
?DOC(
" Check out a resource from the pool, apply the `next` function, then check\n"
" the resource back in.\n"
"\n"
" ```gleam\n"
" let assert Ok(pool) =\n"
" bath.new(fn() { Ok(\"Some pooled resource\") })\n"
" |> bath.start(1000)\n"
"\n"
" use resource <- bath.apply(pool, 1000)\n"
"\n"
" // Do stuff with resource...\n"
"\n"
" // Return the resource to the pool, returning \"Hello!\" to the caller.\n"
" bath.keep()\n"
" |> bath.returning(\"Hello!\")\n"
" ```\n"
).
-spec apply(
gleam@erlang@process:subject(msg(FFZ)),
integer(),
fun((FFZ) -> next(FGC))
) -> {ok, FGC} | {error, apply_error()}.
apply(Pool, Timeout, Next) ->
Self = erlang:self(),
gleam@result:'try'(
check_out(Pool, Self, Timeout),
fun(Resource) ->
Next_action = Next(Resource),
{Usage_result, Next_action@1} = case Next_action of
{keep, Return} ->
{Return, {keep, nil}};
{discard, Return@1} ->
{Return@1, {discard, nil}}
end,
check_in(Pool, Resource, Self, Next_action@1),
{ok, Usage_result}
end
).
-file("src/bath.gleam", 335).
?DOC(
" Like [`apply`](#apply), but blocks waiting for a resource if the pool is\n"
" exhausted instead of returning `NoResourcesAvailable` immediately.\n"
"\n"
" The timeout covers both waiting in the queue and any internal pool operations.\n"
"\n"
" ## Panics\n"
"\n"
" This function will panic if the timeout expires before a resource becomes\n"
" available.\n"
"\n"
" If you need to handle timeouts gracefully, you can:\n"
" - Use this function within a supervised process (let it crash, supervisor restarts)\n"
" - Wrap the call in a separate process/task that you can monitor\n"
"\n"
" If the pool shuts down while waiting, `Error(PoolShuttingDown)` is returned.\n"
"\n"
" ## Example\n"
"\n"
" ```gleam\n"
" let assert Ok(pool) =\n"
" bath.new(fn() { Ok(\"Some pooled resource\") })\n"
" |> bath.size(1)\n"
" |> bath.start(1000)\n"
"\n"
" // This will block until a resource is available or panic on timeout\n"
" use resource <- bath.apply_blocking(pool, 5000)\n"
"\n"
" // Do stuff with resource...\n"
"\n"
" bath.keep()\n"
" |> bath.returning(\"Hello!\")\n"
" ```\n"
).
-spec apply_blocking(
gleam@erlang@process:subject(msg(FGG)),
integer(),
fun((FGG) -> next(FGJ))
) -> {ok, FGJ} | {error, apply_error()}.
apply_blocking(Pool, Timeout, Next) ->
Self = erlang:self(),
gleam@result:'try'(
gleam@erlang@process:call(
Pool,
Timeout,
fun(_capture) -> {check_out_blocking, _capture, Self} end
),
fun(Resource) ->
Next_action = Next(Resource),
{Usage_result, Next_action@1} = case Next_action of
{keep, Return} ->
{Return, {keep, nil}};
{discard, Return@1} ->
{Return@1, {discard, nil}}
end,
check_in(Pool, Resource, Self, Next_action@1),
{ok, Usage_result}
end
).
-file("src/bath.gleam", 367).
?DOC(
" Shut down the pool, calling the shutdown function on each\n"
" resource in the pool. Calling with `force` set to `True` will\n"
" force the shutdown, not calling the shutdown function on any\n"
" resources.\n"
"\n"
" Will fail if there are still resources checked out, unless `force` is\n"
" `True`.\n"
"\n"
" You only need to call this when using unsupervised pools. You should let your\n"
" supervision tree handle the shutdown of supervised resource pools.\n"
).
-spec shutdown(gleam@erlang@process:subject(msg(any())), boolean(), integer()) -> {ok,
nil} |
{error, shutdown_error()}.
shutdown(Pool, Force, Timeout) ->
gleam@erlang@process:call(
Pool,
Timeout,
fun(_capture) -> {shutdown, _capture, Force} end
).
-file("src/bath.gleam", 800).
-spec reject_waiters(gleam@deque:deque(waiting_request(any()))) -> nil.
reject_waiters(Waiting) ->
_pipe = Waiting,
_pipe@1 = gleam@deque:to_list(_pipe),
gleam@list:each(
_pipe@1,
fun(Waiter) ->
gleam@erlang@process:demonitor_process(erlang:element(4, Waiter)),
gleam@otp@actor:send(
erlang:element(2, Waiter),
{error, pool_shutting_down}
)
end
).
-file("src/bath.gleam", 957).
-spec monitor_process(
gleam@erlang@process:selector(msg(FHW)),
gleam@erlang@process:pid_()
) -> {gleam@erlang@process:monitor(), gleam@erlang@process:selector(msg(FHW))}.
monitor_process(Selector, Pid) ->
Monitor = gleam@erlang@process:monitor(Pid),
Selector@1 = begin
_pipe = Selector,
gleam@erlang@process:select_specific_monitor(
_pipe,
Monitor,
fun(Field@0) -> {caller_down, Field@0} end
)
end,
{Monitor, Selector@1}.
-file("src/bath.gleam", 965).
-spec demonitor_process(
gleam@erlang@process:selector(msg(FIA)),
gleam@erlang@process:monitor()
) -> gleam@erlang@process:selector(msg(FIA)).
demonitor_process(Selector, Monitor) ->
gleam@erlang@process:demonitor_process(Monitor),
Selector@1 = begin
_pipe = Selector,
gleam@erlang@process:deselect_specific_monitor(_pipe, Monitor)
end,
Selector@1.
-file("src/bath.gleam", 976).
-spec log_resource_creation_error(boolean(), binary()) -> nil.
log_resource_creation_error(Log_errors, Resource_create_error) ->
case Log_errors of
true ->
logging:log(
error,
<<"Bath: Resource creation failed: "/utf8,
Resource_create_error/binary>>
);
false ->
nil
end.
-file("src/bath.gleam", 910).
?DOC(
" Try to pop the next waiter from the queue and serve them a resource.\n"
"\n"
" If `resource` is `Ok`, that resource is handed to the waiter directly.\n"
" If `resource` is `Error`, a new resource is created for the waiter.\n"
"\n"
" `current_size` should reflect the pool size *before* this resource is\n"
" accounted for (i.e. after the old resource was removed/shutdown).\n"
).
-spec serve_next_waiter(
state(FHO),
{ok, FHO} | {error, nil},
gleam@dict:dict(gleam@erlang@process:pid_(), live_resource(FHO)),
gleam@erlang@process:selector(msg(FHO)),
integer()
) -> serve_result(FHO).
serve_next_waiter(State, Resource, Live_resources, Selector, Current_size) ->
case gleam@deque:pop_front(erlang:element(10, State)) of
{error, _} ->
none_waiting;
{ok, {Waiter, New_waiting}} ->
Resource_result = case Resource of
{ok, R} ->
{ok, {R, Current_size}};
{error, _} ->
case (erlang:element(5, State))() of
{ok, R@1} ->
{ok, {R@1, Current_size + 1}};
{error, Err} ->
{error, Err}
end
end,
case Resource_result of
{ok, {R@2, New_current_size}} ->
Selector@1 = demonitor_process(
Selector,
erlang:element(4, Waiter)
),
{Monitor, Selector@2} = monitor_process(
Selector@1,
erlang:element(3, Waiter)
),
Live_resources@1 = gleam@dict:insert(
Live_resources,
erlang:element(3, Waiter),
{live_resource, R@2, Monitor}
),
gleam@otp@actor:send(erlang:element(2, Waiter), {ok, R@2}),
{served,
Live_resources@1,
New_waiting,
Selector@2,
New_current_size};
{error, Err@1} ->
log_resource_creation_error(
erlang:element(12, State),
Err@1
),
Selector@3 = demonitor_process(
Selector,
erlang:element(4, Waiter)
),
gleam@otp@actor:send(
erlang:element(2, Waiter),
{error, {check_out_resource_create_error, Err@1}}
),
{serve_failed, New_waiting, Selector@3, Current_size}
end
end.
-file("src/bath.gleam", 425).
-spec handle_pool_message(state(FGS), msg(FGS)) -> gleam@otp@actor:next(state(FGS), msg(FGS)).
handle_pool_message(State, Msg) ->
case Msg of
{check_in, Resource, Caller, Next} ->
Caller_live_resource = gleam_stdlib:map_get(
erlang:element(9, State),
Caller
),
Live_resources = gleam@dict:delete(erlang:element(9, State), Caller),
Selector = case Caller_live_resource of
{ok, Live_resource} ->
demonitor_process(
erlang:element(11, State),
erlang:element(3, Live_resource)
);
{error, _} ->
erlang:element(11, State)
end,
case Next of
{keep, _} ->
case serve_next_waiter(
State,
{ok, Resource},
Live_resources,
Selector,
erlang:element(8, State)
) of
{served, Live_resources@1, Waiting, Selector@1, _} ->
_pipe = {state,
erlang:element(2, State),
erlang:element(3, State),
erlang:element(4, State),
erlang:element(5, State),
erlang:element(6, State),
erlang:element(7, State),
erlang:element(8, State),
Live_resources@1,
Waiting,
Selector@1,
erlang:element(12, State)},
_pipe@1 = gleam@otp@actor:continue(_pipe),
gleam@otp@actor:with_selector(_pipe@1, Selector@1);
none_waiting ->
_pipe@2 = {state,
erlang:element(2, State),
erlang:element(3, State),
erlang:element(4, State),
erlang:element(5, State),
erlang:element(6, State),
gleam@deque:push_back(
erlang:element(7, State),
Resource
),
erlang:element(8, State),
Live_resources,
erlang:element(10, State),
Selector,
erlang:element(12, State)},
_pipe@3 = gleam@otp@actor:continue(_pipe@2),
gleam@otp@actor:with_selector(_pipe@3, Selector);
{serve_failed, _, _, _} ->
erlang:error(#{gleam_error => panic,
message => <<"unreachable: serve_next_waiter with Ok resource cannot fail"/utf8>>,
file => <<?FILEPATH/utf8>>,
module => <<"bath"/utf8>>,
function => <<"handle_pool_message"/utf8>>,
line => 466})
end;
{discard, _} ->
(erlang:element(6, State))(Resource),
Current_size = erlang:element(8, State) - 1,
case serve_next_waiter(
State,
{error, nil},
Live_resources,
Selector,
Current_size
) of
{served,
Live_resources@2,
Waiting@1,
Selector@2,
Current_size@1} ->
_pipe@4 = {state,
erlang:element(2, State),
erlang:element(3, State),
erlang:element(4, State),
erlang:element(5, State),
erlang:element(6, State),
erlang:element(7, State),
Current_size@1,
Live_resources@2,
Waiting@1,
Selector@2,
erlang:element(12, State)},
_pipe@5 = gleam@otp@actor:continue(_pipe@4),
gleam@otp@actor:with_selector(_pipe@5, Selector@2);
{serve_failed, Waiting@2, Selector@3, Current_size@2} ->
_pipe@6 = {state,
erlang:element(2, State),
erlang:element(3, State),
erlang:element(4, State),
erlang:element(5, State),
erlang:element(6, State),
erlang:element(7, State),
Current_size@2,
Live_resources,
Waiting@2,
Selector@3,
erlang:element(12, State)},
_pipe@7 = gleam@otp@actor:continue(_pipe@6),
gleam@otp@actor:with_selector(_pipe@7, Selector@3);
none_waiting ->
_pipe@8 = {state,
erlang:element(2, State),
erlang:element(3, State),
erlang:element(4, State),
erlang:element(5, State),
erlang:element(6, State),
erlang:element(7, State),
Current_size,
Live_resources,
erlang:element(10, State),
Selector,
erlang:element(12, State)},
_pipe@9 = gleam@otp@actor:continue(_pipe@8),
gleam@otp@actor:with_selector(_pipe@9, Selector)
end
end;
{check_out, Reply_to, Caller@1} ->
Get_result = case erlang:element(2, State) of
f_i_f_o ->
gleam@deque:pop_front(erlang:element(7, State));
l_i_f_o ->
gleam@deque:pop_back(erlang:element(7, State))
end,
Resource_result = case Get_result of
{ok, {Resource@1, New_resources}} ->
{ok, {Resource@1, New_resources, erlang:element(8, State)}};
{error, _} ->
case erlang:element(8, State) < erlang:element(4, State) of
true ->
gleam@result:'try'(
begin
_pipe@10 = (erlang:element(5, State))(),
gleam@result:map_error(
_pipe@10,
fun(Err) ->
log_resource_creation_error(
erlang:element(12, State),
Err
),
{check_out_resource_create_error,
Err}
end
)
end,
fun(Resource@2) ->
{ok,
{Resource@2,
erlang:element(7, State),
erlang:element(8, State) + 1}}
end
);
false ->
{error, no_resources_available}
end
end,
case Resource_result of
{error, Err@1} ->
gleam@otp@actor:send(Reply_to, {error, Err@1}),
gleam@otp@actor:continue(State);
{ok, {Resource@3, New_resources@1, New_current_size}} ->
{Monitor, Selector@4} = monitor_process(
erlang:element(11, State),
Caller@1
),
Live_resources@3 = gleam@dict:insert(
erlang:element(9, State),
Caller@1,
{live_resource, Resource@3, Monitor}
),
gleam@otp@actor:send(Reply_to, {ok, Resource@3}),
_pipe@11 = {state,
erlang:element(2, State),
erlang:element(3, State),
erlang:element(4, State),
erlang:element(5, State),
erlang:element(6, State),
New_resources@1,
New_current_size,
Live_resources@3,
erlang:element(10, State),
Selector@4,
erlang:element(12, State)},
_pipe@12 = gleam@otp@actor:continue(_pipe@11),
gleam@otp@actor:with_selector(_pipe@12, Selector@4)
end;
{check_out_blocking, Reply_to@1, Caller@2} ->
Get_result@1 = case erlang:element(2, State) of
f_i_f_o ->
gleam@deque:pop_front(erlang:element(7, State));
l_i_f_o ->
gleam@deque:pop_back(erlang:element(7, State))
end,
Resource_result@1 = case Get_result@1 of
{ok, {Resource@4, New_resources@2}} ->
{ok,
{Resource@4, New_resources@2, erlang:element(8, State)}};
{error, _} ->
case erlang:element(8, State) < erlang:element(4, State) of
true ->
gleam@result:'try'(
begin
_pipe@13 = (erlang:element(5, State))(),
gleam@result:map_error(
_pipe@13,
fun(Err@2) ->
log_resource_creation_error(
erlang:element(12, State),
Err@2
),
{check_out_resource_create_error,
Err@2}
end
)
end,
fun(Resource@5) ->
{ok,
{Resource@5,
erlang:element(7, State),
erlang:element(8, State) + 1}}
end
);
false ->
{error, no_resources_available}
end
end,
case Resource_result@1 of
{error, no_resources_available} ->
Monitor@1 = gleam@erlang@process:monitor(Caller@2),
Selector@5 = begin
_pipe@14 = erlang:element(11, State),
gleam@erlang@process:select_specific_monitor(
_pipe@14,
Monitor@1,
fun(Field@0) -> {waiter_down, Field@0} end
)
end,
Waiting_request = {waiting_request,
Reply_to@1,
Caller@2,
Monitor@1},
_pipe@15 = {state,
erlang:element(2, State),
erlang:element(3, State),
erlang:element(4, State),
erlang:element(5, State),
erlang:element(6, State),
erlang:element(7, State),
erlang:element(8, State),
erlang:element(9, State),
gleam@deque:push_back(
erlang:element(10, State),
Waiting_request
),
Selector@5,
erlang:element(12, State)},
_pipe@16 = gleam@otp@actor:continue(_pipe@15),
gleam@otp@actor:with_selector(_pipe@16, Selector@5);
{error, Err@3} ->
gleam@otp@actor:send(Reply_to@1, {error, Err@3}),
gleam@otp@actor:continue(State);
{ok, {Resource@6, New_resources@3, New_current_size@1}} ->
{Monitor@2, Selector@6} = monitor_process(
erlang:element(11, State),
Caller@2
),
Live_resources@4 = gleam@dict:insert(
erlang:element(9, State),
Caller@2,
{live_resource, Resource@6, Monitor@2}
),
gleam@otp@actor:send(Reply_to@1, {ok, Resource@6}),
_pipe@17 = {state,
erlang:element(2, State),
erlang:element(3, State),
erlang:element(4, State),
erlang:element(5, State),
erlang:element(6, State),
New_resources@3,
New_current_size@1,
Live_resources@4,
erlang:element(10, State),
Selector@6,
erlang:element(12, State)},
_pipe@18 = gleam@otp@actor:continue(_pipe@17),
gleam@otp@actor:with_selector(_pipe@18, Selector@6)
end;
{pool_exit, Exit_message} ->
_pipe@19 = erlang:element(7, State),
_pipe@20 = gleam@deque:to_list(_pipe@19),
gleam@list:each(_pipe@20, erlang:element(6, State)),
case erlang:element(3, Exit_message) of
{abnormal, Reason} ->
_pipe@21 = gleam@string:inspect(Reason),
gleam@otp@actor:stop_abnormal(_pipe@21);
killed ->
gleam@otp@actor:stop_abnormal(<<"Killed"/utf8>>);
normal ->
gleam@otp@actor:stop()
end;
{shutdown, Reply_to@2, Force} ->
case {maps:size(erlang:element(9, State)), Force} of
{0, _} ->
reject_waiters(erlang:element(10, State)),
_pipe@22 = erlang:element(7, State),
_pipe@23 = gleam@deque:to_list(_pipe@22),
gleam@list:each(_pipe@23, erlang:element(6, State)),
gleam@otp@actor:send(Reply_to@2, {ok, nil}),
gleam@otp@actor:stop();
{_, true} ->
reject_waiters(erlang:element(10, State)),
gleam@otp@actor:send(Reply_to@2, {ok, nil}),
gleam@otp@actor:stop();
{_, false} ->
gleam@otp@actor:send(Reply_to@2, {error, resources_in_use}),
gleam@otp@actor:continue(State)
end;
{caller_down, Process_down} ->
Process_down_pid@1 = case Process_down of
{process_down, _, Process_down_pid, _} -> Process_down_pid;
_assert_fail ->
erlang:error(#{gleam_error => let_assert,
message => <<"Pattern match failed, no pattern matched the value."/utf8>>,
file => <<?FILEPATH/utf8>>,
module => <<"bath"/utf8>>,
function => <<"handle_pool_message"/utf8>>,
line => 683,
value => _assert_fail,
start => 20584,
'end' => 20656,
pattern_start => 20595,
pattern_end => 20641})
end,
case gleam_stdlib:map_get(
erlang:element(9, State),
Process_down_pid@1
) of
{error, _} ->
gleam@otp@actor:continue(State);
{ok, Live_resource@1} ->
Selector@7 = demonitor_process(
erlang:element(11, State),
erlang:element(3, Live_resource@1)
),
Live_resources@5 = gleam@dict:delete(
erlang:element(9, State),
Process_down_pid@1
),
(erlang:element(6, State))(
erlang:element(2, Live_resource@1)
),
Current_size@3 = erlang:element(8, State),
case serve_next_waiter(
State,
{error, nil},
Live_resources@5,
Selector@7,
Current_size@3 - 1
) of
{served,
Live_resources@6,
Waiting@3,
Selector@8,
Current_size@4} ->
_pipe@24 = {state,
erlang:element(2, State),
erlang:element(3, State),
erlang:element(4, State),
erlang:element(5, State),
erlang:element(6, State),
erlang:element(7, State),
Current_size@4,
Live_resources@6,
Waiting@3,
Selector@8,
erlang:element(12, State)},
_pipe@25 = gleam@otp@actor:continue(_pipe@24),
gleam@otp@actor:with_selector(_pipe@25, Selector@8);
{serve_failed, Waiting@4, Selector@9, Current_size@5} ->
_pipe@26 = {state,
erlang:element(2, State),
erlang:element(3, State),
erlang:element(4, State),
erlang:element(5, State),
erlang:element(6, State),
erlang:element(7, State),
Current_size@5,
Live_resources@5,
Waiting@4,
Selector@9,
erlang:element(12, State)},
_pipe@27 = gleam@otp@actor:continue(_pipe@26),
gleam@otp@actor:with_selector(_pipe@27, Selector@9);
none_waiting ->
{New_resources@4, New_current_size@2} = case erlang:element(
3,
State
) of
lazy ->
{erlang:element(7, State),
Current_size@3 - 1};
eager ->
case (erlang:element(5, State))() of
{ok, Resource@7} ->
{gleam@deque:push_back(
erlang:element(7, State),
Resource@7
),
Current_size@3};
{error, Resource_create_error} ->
log_resource_creation_error(
erlang:element(12, State),
Resource_create_error
),
{erlang:element(7, State),
Current_size@3 - 1}
end
end,
_pipe@28 = {state,
erlang:element(2, State),
erlang:element(3, State),
erlang:element(4, State),
erlang:element(5, State),
erlang:element(6, State),
New_resources@4,
New_current_size@2,
Live_resources@5,
erlang:element(10, State),
Selector@7,
erlang:element(12, State)},
_pipe@29 = gleam@otp@actor:continue(_pipe@28),
gleam@otp@actor:with_selector(_pipe@29, Selector@7)
end
end;
{waiter_down, Process_down@1} ->
Process_down_pid@3 = case Process_down@1 of
{process_down, _, Process_down_pid@2, _} -> Process_down_pid@2;
_assert_fail@1 ->
erlang:error(#{gleam_error => let_assert,
message => <<"Pattern match failed, no pattern matched the value."/utf8>>,
file => <<?FILEPATH/utf8>>,
module => <<"bath"/utf8>>,
function => <<"handle_pool_message"/utf8>>,
line => 777,
value => _assert_fail@1,
start => 24147,
'end' => 24219,
pattern_start => 24158,
pattern_end => 24204})
end,
{New_waiting, Monitors_to_remove} = begin
_pipe@30 = erlang:element(10, State),
_pipe@31 = gleam@deque:to_list(_pipe@30),
_pipe@32 = gleam@list:partition(
_pipe@31,
fun(Waiter) ->
erlang:element(3, Waiter) /= Process_down_pid@3
end
),
(fun(Partitioned) ->
{Kept, Removed} = Partitioned,
{gleam@deque:from_list(Kept),
gleam@list:map(
Removed,
fun(W) -> erlang:element(4, W) end
)}
end)(_pipe@32)
end,
Selector@10 = gleam@list:fold(
Monitors_to_remove,
erlang:element(11, State),
fun(Sel, Mon) -> demonitor_process(Sel, Mon) end
),
_pipe@33 = {state,
erlang:element(2, State),
erlang:element(3, State),
erlang:element(4, State),
erlang:element(5, State),
erlang:element(6, State),
erlang:element(7, State),
erlang:element(8, State),
erlang:element(9, State),
New_waiting,
Selector@10,
erlang:element(12, State)},
_pipe@34 = gleam@otp@actor:continue(_pipe@33),
gleam@otp@actor:with_selector(_pipe@34, Selector@10)
end.
-file("src/bath.gleam", 994).
?DOC(false).
-spec try_map_returning(list(FIE), fun((FIE) -> {ok, FIG} | {error, FIH})) -> {ok,
list(FIG)} |
{error, {list(FIG), FIH}}.
try_map_returning(List, Fun) ->
_pipe = List,
_pipe@1 = gleam@list:try_fold(_pipe, [], fun(Acc, Item) -> case Fun(Item) of
{ok, Value} ->
{ok, [Value | Acc]};
{error, Error} ->
{error, {Acc, Error}}
end end),
gleam@result:map(_pipe@1, fun lists:reverse/1).
-file("src/bath.gleam", 811).
?DOC(
" Create the resources for a pool, returning a deque of resources and the number of\n"
" resources created.\n"
).
-spec create_pool_resources(builder(FGZ)) -> {ok,
{gleam@deque:deque(FGZ), integer()}} |
{error, binary()}.
create_pool_resources(Builder) ->
case erlang:element(7, Builder) of
lazy ->
{ok, {gleam@deque:new(), 0}};
eager ->
Create_result = begin
_pipe = gleam@list:repeat(
<<""/utf8>>,
erlang:element(3, Builder)
),
_pipe@1 = try_map_returning(
_pipe,
fun(_) -> (erlang:element(4, Builder))() end
),
gleam@result:map(_pipe@1, fun gleam@deque:from_list/1)
end,
case Create_result of
{ok, Resources} ->
{ok, {Resources, erlang:element(3, Builder)}};
{error, {Created_resources, Error}} ->
_pipe@2 = Created_resources,
gleam@list:each(_pipe@2, erlang:element(5, Builder)),
{error, Error}
end
end.
-file("src/bath.gleam", 835).
-spec actor_builder(
builder(FHE),
fun((gleam@erlang@process:subject(msg(FHE))) -> FHI),
integer()
) -> gleam@otp@actor:builder(state(FHE), msg(FHE), FHI).
actor_builder(Builder, Mapper, Init_timeout) ->
Pool_builder = begin
_pipe@5 = gleam@otp@actor:new_with_initialiser(
Init_timeout,
fun(Self) ->
gleam@result:'try'(
create_pool_resources(Builder),
fun(_use0) ->
{Resources, Current_size} = _use0,
gleam_erlang_ffi:trap_exits(true),
Selector = begin
_pipe = gleam_erlang_ffi:new_selector(),
_pipe@1 = gleam@erlang@process:select(_pipe, Self),
gleam@erlang@process:select_trapped_exits(
_pipe@1,
fun(Field@0) -> {pool_exit, Field@0} end
)
end,
State = {state,
erlang:element(6, Builder),
erlang:element(7, Builder),
erlang:element(3, Builder),
erlang:element(4, Builder),
erlang:element(5, Builder),
Resources,
Current_size,
maps:new(),
gleam@deque:new(),
Selector,
erlang:element(8, Builder)},
_pipe@2 = gleam@otp@actor:initialised(State),
_pipe@3 = gleam@otp@actor:selecting(_pipe@2, Selector),
_pipe@4 = gleam@otp@actor:returning(
_pipe@3,
Mapper(Self)
),
{ok, _pipe@4}
end
)
end
),
gleam@otp@actor:on_message(_pipe@5, fun handle_pool_message/2)
end,
case erlang:element(2, Builder) of
{some, Name} ->
gleam@otp@actor:named(Pool_builder, Name);
none ->
Pool_builder
end.
-file("src/bath.gleam", 191).
?DOC(
" Like [`supervised`](#supervised), but allows you to pass a mapping function to\n"
" transform the pool return value to the receiver. This is mostly useful for library\n"
" authors who wish to use Bath to create a pool of resources.\n"
).
-spec supervised_map(
builder(FEX),
fun((gleam@erlang@process:subject(msg(FEX))) -> FFB),
integer()
) -> gleam@otp@supervision:child_specification(FFB).
supervised_map(Builder, Mapper, Init_timeout) ->
gleam@otp@supervision:worker(
fun() -> _pipe = actor_builder(Builder, Mapper, Init_timeout),
gleam@otp@actor:start(_pipe) end
).
-file("src/bath.gleam", 181).
?DOC(
" Return the [`ChildSpecification`](https://hexdocs.pm/gleam_otp/gleam/otp/supervision.html#ChildSpecification)\n"
" for creating a supervised resource pool.\n"
"\n"
" In order to use a supervised pool, your pool _must_ be named, otherwise you will\n"
" not be able to send messages to your pool. See the [`name`](#name) function.\n"
"\n"
" ## Example\n"
"\n"
" ```gleam\n"
" import bath\n"
" import gleam/erlang/process\n"
" import gleam/otp/static_supervisor as supervisor\n"
"\n"
" fn main() {\n"
" // Create a name to interact with the pool once it's started under the\n"
" // static supervisor.\n"
" let pool_name = process.new_name(\"bath_pool\")\n"
"\n"
" let assert Ok(_started) =\n"
" supervisor.new(supervisor.OneForOne)\n"
" |> supervisor.add(\n"
" bath.new(create_resource)\n"
" |> bath.name(pool_name)\n"
" |> bath.supervised(1000)\n"
" )\n"
" |> supervisor.start\n"
"\n"
" let pool = process.named_subject(pool_name)\n"
"\n"
" // Do more stuff...\n"
" }\n"
" ```\n"
).
-spec supervised(builder(FEU), integer()) -> gleam@otp@supervision:child_specification(gleam@erlang@process:subject(msg(FEU))).
supervised(Builder, Init_timeout) ->
supervised_map(Builder, fun gleam@function:identity/1, Init_timeout).
-file("src/bath.gleam", 205).
?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"
).
-spec start(builder(FFD), integer()) -> {ok,
gleam@erlang@process:subject(msg(FFD))} |
{error, gleam@otp@actor:start_error()}.
start(Builder, Init_timeout) ->
gleam@result:'try'(
begin
_pipe = actor_builder(
Builder,
fun gleam@function:identity/1,
Init_timeout
),
gleam@otp@actor:start(_pipe)
end,
fun(Started) -> {ok, erlang:element(3, Started)} end
).