Current section

Files

Jump to
bath src bath.erl
Raw

src/bath.erl

-module(bath).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch]).
-define(FILEPATH, "src/bath.gleam").
-export([new/1, size/2, on_shutdown/2, checkout_strategy/2, creation_strategy/2, log_errors/2, keep/0, discard/0, returning/2, apply/3, shutdown/3, try_map_returning/2, supervised_map/4, supervised/3, start/2]).
-export_type([checkout_strategy/0, creation_strategy/0, builder/1, apply_error/0, shutdown_error/0, next/1, pool/1, state/1, live_resource/1, msg/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(FJQ) :: {builder,
integer(),
fun(() -> {ok, FJQ} | {error, binary()}),
fun((FJQ) -> nil),
checkout_strategy(),
creation_strategy(),
boolean()}.
-type apply_error() :: no_resources_available |
{check_out_resource_create_error, binary()}.
-type shutdown_error() :: resources_in_use.
-opaque next(FJR) :: {keep, FJR} | {discard, FJR}.
-opaque pool(FJS) :: {pool, gleam@erlang@process:subject(msg(FJS))}.
-opaque state(FJT) :: {state,
checkout_strategy(),
creation_strategy(),
integer(),
fun(() -> {ok, FJT} | {error, binary()}),
fun((FJT) -> nil),
gleam@deque:deque(FJT),
integer(),
gleam@dict:dict(gleam@erlang@process:pid_(), live_resource(FJT)),
gleam@erlang@process:selector(msg(FJT)),
boolean()}.
-type live_resource(FJU) :: {live_resource, FJU, gleam@erlang@process:monitor()}.
-opaque msg(FJV) :: {check_in, FJV, gleam@erlang@process:pid_(), next(nil)} |
{check_out,
gleam@erlang@process:subject({ok, FJV} | {error, apply_error()}),
gleam@erlang@process:pid_()} |
{pool_exit, gleam@erlang@process:exit_message()} |
{caller_down, gleam@erlang@process:down()} |
{shutdown,
gleam@erlang@process:subject({ok, nil} | {error, shutdown_error()}),
boolean()}.
-file("src/bath.gleam", 65).
?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.with_size(10)\n"
" |> bath.start(1000)\n"
" }\n"
" ```\n"
"\n"
" ### Default values\n"
"\n"
" | Config | Default |\n"
" |--------|---------|\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, FKA} | {error, binary()})) -> builder(FKA).
new(Create_resource) ->
{builder, 10, Create_resource, fun(_) -> nil end, f_i_f_o, lazy, false}.
-file("src/bath.gleam", 79).
?DOC(" Set the pool size. Defaults to 10. Will be clamped to a minimum of 1.\n").
-spec size(builder(FKE), integer()) -> builder(FKE).
size(Builder, Size) ->
_record = Builder,
{builder,
gleam@int:max(Size, 1),
erlang:element(3, _record),
erlang:element(4, _record),
erlang:element(5, _record),
erlang:element(6, _record),
erlang:element(7, _record)}.
-file("src/bath.gleam", 87).
?DOC(" Set a shutdown function to be run for each resource when the pool exits.\n").
-spec on_shutdown(builder(FKH), fun((FKH) -> nil)) -> builder(FKH).
on_shutdown(Builder, Shutdown_resource) ->
_record = Builder,
{builder,
erlang:element(2, _record),
erlang:element(3, _record),
Shutdown_resource,
erlang:element(5, _record),
erlang:element(6, _record),
erlang:element(7, _record)}.
-file("src/bath.gleam", 95).
?DOC(" Change the checkout strategy for the pool. Defaults to `FIFO`.\n").
-spec checkout_strategy(builder(FKK), checkout_strategy()) -> builder(FKK).
checkout_strategy(Builder, Checkout_strategy) ->
_record = Builder,
{builder,
erlang:element(2, _record),
erlang:element(3, _record),
erlang:element(4, _record),
Checkout_strategy,
erlang:element(6, _record),
erlang:element(7, _record)}.
-file("src/bath.gleam", 103).
?DOC(" Change the resource creation strategy for the pool. Defaults to `Lazy`.\n").
-spec creation_strategy(builder(FKN), creation_strategy()) -> builder(FKN).
creation_strategy(Builder, Creation_strategy) ->
_record = Builder,
{builder,
erlang:element(2, _record),
erlang:element(3, _record),
erlang:element(4, _record),
erlang:element(5, _record),
Creation_strategy,
erlang:element(7, _record)}.
-file("src/bath.gleam", 111).
?DOC(" Set whether the pool logs errors when resources fail to create.\n").
-spec log_errors(builder(FKQ), boolean()) -> builder(FKQ).
log_errors(Builder, Log_errors) ->
_record = Builder,
{builder,
erlang:element(2, _record),
erlang:element(3, _record),
erlang:element(4, _record),
erlang:element(5, _record),
erlang:element(6, _record),
Log_errors}.
-file("src/bath.gleam", 215).
?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", 224).
?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", 229).
?DOC(" Return a value from a use of [`apply`](#apply).\n").
-spec returning(next(any()), FLN) -> next(FLN).
returning(Next, Value) ->
case Next of
{keep, _} ->
{keep, Value};
{discard, _} ->
{discard, Value}
end.
-file("src/bath.gleam", 239).
?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(pool(FLP), gleam@erlang@process:pid_(), integer()) -> {ok, FLP} |
{error, apply_error()}.
check_out(Pool, Caller, Timeout) ->
gleam@erlang@process:call(
erlang:element(2, Pool),
Timeout,
fun(_capture) -> {check_out, _capture, Caller} end
).
-file("src/bath.gleam", 247).
-spec check_in(pool(FLT), FLT, gleam@erlang@process:pid_(), next(nil)) -> nil.
check_in(Pool, Resource, Caller, Next) ->
gleam@erlang@process:send(
erlang:element(2, Pool),
{check_in, Resource, Caller, Next}
).
-file("src/bath.gleam", 272).
?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(pool(FLX), integer(), fun((FLX) -> next(FLZ))) -> {ok, FLZ} |
{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", 299).
?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"
).
-spec shutdown(pool(any()), boolean(), integer()) -> {ok, nil} |
{error, shutdown_error()}.
shutdown(Pool, Force, Timeout) ->
gleam@erlang@process:call(
erlang:element(2, Pool),
Timeout,
fun(_capture) -> {shutdown, _capture, Force} end
).
-file("src/bath.gleam", 611).
-spec monitor_process(
gleam@erlang@process:selector(msg(FMZ)),
gleam@erlang@process:pid_()
) -> {gleam@erlang@process:monitor(), gleam@erlang@process:selector(msg(FMZ))}.
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", 619).
-spec demonitor_process(
gleam@erlang@process:selector(msg(FND)),
gleam@erlang@process:monitor()
) -> gleam@erlang@process:selector(msg(FND)).
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", 630).
-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", 348).
-spec handle_pool_message(state(FMH), msg(FMH)) -> gleam@otp@actor:next(state(FMH), msg(FMH)).
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(10, State),
erlang:element(3, Live_resource)
);
{error, _} ->
erlang:element(10, State)
end,
{New_resources, Current_size} = case Next of
{keep, _} ->
{gleam@deque:push_back(erlang:element(7, State), Resource),
erlang:element(8, State)};
{discard, _} ->
(erlang:element(6, State))(Resource),
{erlang:element(7, State), erlang:element(8, State) - 1}
end,
_pipe = begin
_record = State,
{state,
erlang:element(2, _record),
erlang:element(3, _record),
erlang:element(4, _record),
erlang:element(5, _record),
erlang:element(6, _record),
New_resources,
Current_size,
Live_resources,
Selector,
erlang:element(11, _record)}
end,
_pipe@1 = gleam@otp@actor:continue(_pipe),
gleam@otp@actor:with_selector(_pipe@1, Selector);
{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@1}} ->
{ok,
{Resource@1, New_resources@1, erlang:element(8, State)}};
{error, _} ->
case erlang:element(8, State) < erlang:element(4, State) of
true ->
gleam@result:'try'(
begin
_pipe@2 = (erlang:element(5, State))(),
gleam@result:map_error(
_pipe@2,
fun(Err) ->
log_resource_creation_error(
erlang:element(11, 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@2, New_current_size}} ->
{Monitor, Selector@1} = monitor_process(
erlang:element(10, State),
Caller@1
),
Live_resources@1 = gleam@dict:insert(
erlang:element(9, State),
Caller@1,
{live_resource, Resource@3, Monitor}
),
gleam@otp@actor:send(Reply_to, {ok, Resource@3}),
_pipe@3 = begin
_record@1 = State,
{state,
erlang:element(2, _record@1),
erlang:element(3, _record@1),
erlang:element(4, _record@1),
erlang:element(5, _record@1),
erlang:element(6, _record@1),
New_resources@2,
New_current_size,
Live_resources@1,
Selector@1,
erlang:element(11, _record@1)}
end,
_pipe@4 = gleam@otp@actor:continue(_pipe@3),
gleam@otp@actor:with_selector(_pipe@4, Selector@1)
end;
{pool_exit, Exit_message} ->
_pipe@5 = erlang:element(7, State),
_pipe@6 = gleam@deque:to_list(_pipe@5),
gleam@list:each(_pipe@6, erlang:element(6, State)),
case erlang:element(3, Exit_message) of
{abnormal, Reason} ->
_pipe@7 = gleam@string:inspect(Reason),
gleam@otp@actor:stop_abnormal(_pipe@7);
killed ->
gleam@otp@actor:stop_abnormal(<<"Killed"/utf8>>);
normal ->
gleam@otp@actor:stop()
end;
{shutdown, Reply_to@1, Force} ->
case {maps:size(erlang:element(9, State)), Force} of
{0, _} ->
_pipe@8 = erlang:element(7, State),
_pipe@9 = gleam@deque:to_list(_pipe@8),
gleam@list:each(_pipe@9, erlang:element(6, State)),
gleam@otp@actor:send(Reply_to@1, {ok, nil}),
gleam@otp@actor:stop();
{_, true} ->
gleam@otp@actor:send(Reply_to@1, {ok, nil}),
gleam@otp@actor:stop();
{_, false} ->
gleam@otp@actor:send(Reply_to@1, {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 => 486,
value => _assert_fail,
start => 14227,
'end' => 14299,
pattern_start => 14238,
pattern_end => 14284})
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@2 = demonitor_process(
erlang:element(10, State),
erlang:element(3, Live_resource@1)
),
(erlang:element(6, State))(
erlang:element(2, Live_resource@1)
),
{New_resources@3, New_current_size@1} = case erlang:element(
3,
State
) of
lazy ->
{erlang:element(7, State),
erlang:element(8, State) - 1};
eager ->
case (erlang:element(5, State))() of
{ok, Resource@4} ->
{gleam@deque:push_back(
erlang:element(7, State),
Resource@4
),
erlang:element(8, State)};
{error, Resource_create_error} ->
log_resource_creation_error(
erlang:element(11, State),
Resource_create_error
),
{erlang:element(7, State),
erlang:element(8, State)}
end
end,
_pipe@10 = begin
_record@2 = State,
{state,
erlang:element(2, _record@2),
erlang:element(3, _record@2),
erlang:element(4, _record@2),
erlang:element(5, _record@2),
erlang:element(6, _record@2),
New_resources@3,
New_current_size@1,
gleam@dict:delete(
erlang:element(9, State),
Process_down_pid@1
),
Selector@2,
erlang:element(11, _record@2)}
end,
_pipe@11 = gleam@otp@actor:continue(_pipe@10),
gleam@otp@actor:with_selector(_pipe@11, Selector@2)
end
end.
-file("src/bath.gleam", 648).
?DOC(false).
-spec try_map_returning(list(FNH), fun((FNH) -> {ok, FNJ} | {error, FNK})) -> {ok,
list(FNJ)} |
{error, {list(FNJ), FNK}}.
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", 544).
?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(FML)) -> {ok,
{gleam@deque:deque(FML), integer()}} |
{error, binary()}.
create_pool_resources(Builder) ->
case erlang:element(6, Builder) of
lazy ->
{ok, {gleam@deque:new(), 0}};
eager ->
Create_result = begin
_pipe = gleam@list:repeat(
<<""/utf8>>,
erlang:element(2, Builder)
),
_pipe@1 = try_map_returning(
_pipe,
fun(_) -> (erlang:element(3, Builder))() end
),
gleam@result:map(_pipe@1, fun gleam@deque:from_list/1)
end,
case Create_result of
{ok, Resources} ->
{ok, {Resources, erlang:element(2, Builder)}};
{error, {Created_resources, Error}} ->
_pipe@2 = Created_resources,
gleam@list:each(_pipe@2, erlang:element(4, Builder)),
{error, Error}
end
end.
-file("src/bath.gleam", 568).
-spec actor_builder(builder(FMQ), integer()) -> gleam@otp@actor:builder(state(FMQ), msg(FMQ), gleam@erlang@process:subject(msg(FMQ))).
actor_builder(Builder, Init_timeout) ->
_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(5, Builder),
erlang:element(6, Builder),
erlang:element(2, Builder),
erlang:element(3, Builder),
erlang:element(4, Builder),
Resources,
Current_size,
maps:new(),
Selector,
erlang:element(7, 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, Self),
{ok, _pipe@4}
end
)
end
),
gleam@otp@actor:on_message(_pipe@5, fun handle_pool_message/2).
-file("src/bath.gleam", 174).
?DOC(
" Like [`supervised`](#supervised), but allows you to pass a mapping function to\n"
" transform the pool handler before sending it to the receiver. This is mostly\n"
" useful for library authors who wish to use Bath to create a pool of resources.\n"
).
-spec supervised_map(
builder(FKY),
gleam@erlang@process:subject(FLA),
fun((pool(FKY)) -> FLA),
integer()
) -> gleam@otp@supervision:child_specification(gleam@erlang@process:subject(msg(FKY))).
supervised_map(Builder, Pool_receiver, Mapper, Init_timeout) ->
gleam@otp@supervision:worker(
fun() ->
gleam@result:'try'(
begin
_pipe = actor_builder(Builder, Init_timeout),
gleam@otp@actor:start(_pipe)
end,
fun(Started) ->
gleam@erlang@process:send(
Pool_receiver,
Mapper({pool, erlang:element(3, Started)})
),
{ok, Started}
end
)
end
).
-file("src/bath.gleam", 163).
?DOC(
" Return the [`ChildSpecification`](https://hexdocs.pm/gleam_otp/gleam/otp/supervision.html#ChildSpecification)\n"
" for creating a supervised resource 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 bath\n"
" import gleam/erlang/process\n"
" import gleam/otp/static_supervisor as supervisor\n"
"\n"
" fn main() {\n"
" let pool_receiver = process.new_subject()\n"
"\n"
" let assert Ok(_started) =\n"
" supervisor.new(supervisor.OneForOne)\n"
" |> supervisor.add(\n"
" bath.new(create_resource)\n"
" |> bath.supervised(pool_receiver, 1000)\n"
" )\n"
" |> supervisor.start\n"
"\n"
" let assert Ok(pool) =\n"
" process.receive(pool_receiver)\n"
"\n"
" let assert Ok(_) = bath.apply(pool, fn(res) { echo res })\n"
" }\n"
" ```\n"
).
-spec supervised(
builder(FKT),
gleam@erlang@process:subject(pool(FKT)),
integer()
) -> gleam@otp@supervision:child_specification(gleam@erlang@process:subject(msg(FKT))).
supervised(Builder, Pool_receiver, Init_timeout) ->
supervised_map(
Builder,
Pool_receiver,
fun gleam@function:identity/1,
Init_timeout
).
-file("src/bath.gleam", 194).
?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(FLE), integer()) -> {ok, pool(FLE)} |
{error, gleam@otp@actor:start_error()}.
start(Builder, Init_timeout) ->
gleam@result:'try'(
begin
_pipe = actor_builder(Builder, Init_timeout),
gleam@otp@actor:start(_pipe)
end,
fun(Started) -> {ok, {pool, erlang:element(3, Started)}} end
).