-module(bath). -compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch]). -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, 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, 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, gleam@option:option(gleam@erlang@process:name(msg(FJQ))), 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 state(FJS) :: {state, checkout_strategy(), creation_strategy(), integer(), fun(() -> {ok, FJS} | {error, binary()}), fun((FJS) -> nil), gleam@deque:deque(FJS), integer(), gleam@dict:dict(gleam@erlang@process:pid_(), live_resource(FJS)), gleam@erlang@process:selector(msg(FJS)), boolean()}. -type live_resource(FJT) :: {live_resource, FJT, gleam@erlang@process:monitor()}. -opaque msg(FJU) :: {check_in, FJU, gleam@erlang@process:pid_(), next(nil)} | {check_out, gleam@erlang@process:subject({ok, FJU} | {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", 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.with_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, FJZ} | {error, binary()})) -> builder(FJZ). 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(FKD), integer()) -> builder(FKD). size(Builder, Size) -> _record = Builder, {builder, erlang:element(2, _record), gleam@int:max(Size, 1), erlang:element(4, _record), erlang:element(5, _record), erlang:element(6, _record), erlang:element(7, _record), erlang:element(8, _record)}. -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(FKG), gleam@erlang@process:name(msg(FKG))) -> builder(FKG). name(Builder, Name) -> _record = Builder, {builder, {some, Name}, erlang:element(3, _record), erlang:element(4, _record), erlang:element(5, _record), erlang:element(6, _record), erlang:element(7, _record), erlang:element(8, _record)}. -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(FKL), fun((FKL) -> nil)) -> builder(FKL). on_shutdown(Builder, Shutdown_resource) -> _record = Builder, {builder, erlang:element(2, _record), erlang:element(3, _record), erlang:element(4, _record), Shutdown_resource, erlang:element(6, _record), erlang:element(7, _record), erlang:element(8, _record)}. -file("src/bath.gleam", 110). ?DOC(" Change the checkout strategy for the pool. Defaults to `FIFO`.\n"). -spec checkout_strategy(builder(FKO), checkout_strategy()) -> builder(FKO). checkout_strategy(Builder, Checkout_strategy) -> _record = Builder, {builder, erlang:element(2, _record), erlang:element(3, _record), erlang:element(4, _record), erlang:element(5, _record), Checkout_strategy, erlang:element(7, _record), erlang:element(8, _record)}. -file("src/bath.gleam", 118). ?DOC(" Change the resource creation strategy for the pool. Defaults to `Lazy`.\n"). -spec creation_strategy(builder(FKR), creation_strategy()) -> builder(FKR). creation_strategy(Builder, Creation_strategy) -> _record = Builder, {builder, erlang:element(2, _record), erlang:element(3, _record), erlang:element(4, _record), erlang:element(5, _record), erlang:element(6, _record), Creation_strategy, erlang:element(8, _record)}. -file("src/bath.gleam", 126). ?DOC(" Set whether the pool logs errors when resources fail to create.\n"). -spec log_errors(builder(FKU), boolean()) -> builder(FKU). 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), erlang:element(7, _record), Log_errors}. -file("src/bath.gleam", 225). ?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", 234). ?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", 239). ?DOC(" Return a value from a use of [`apply`](#apply).\n"). -spec returning(next(any()), FLQ) -> next(FLQ). returning(Next, Value) -> case Next of {keep, _} -> {keep, Value}; {discard, _} -> {discard, Value} end. -file("src/bath.gleam", 249). ?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(FLS)), gleam@erlang@process:pid_(), integer() ) -> {ok, FLS} | {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", 257). -spec check_in( gleam@erlang@process:subject(msg(FLX)), FLX, 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", 282). ?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(FMC)), integer(), fun((FMC) -> next(FMF)) ) -> {ok, FMF} | {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", 312). ?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", 624). -spec monitor_process( gleam@erlang@process:selector(msg(FNH)), gleam@erlang@process:pid_() ) -> {gleam@erlang@process:monitor(), gleam@erlang@process:selector(msg(FNH))}. 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", 632). -spec demonitor_process( gleam@erlang@process:selector(msg(FNL)), gleam@erlang@process:monitor() ) -> gleam@erlang@process:selector(msg(FNL)). 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", 643). -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", 356). -spec handle_pool_message(state(FMO), msg(FMO)) -> gleam@otp@actor:next(state(FMO), msg(FMO)). 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 => <>, module => <<"bath"/utf8>>, function => <<"handle_pool_message"/utf8>>, line => 494, value => _assert_fail, start => 14682, 'end' => 14754, pattern_start => 14693, pattern_end => 14739}) 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", 661). ?DOC(false). -spec try_map_returning(list(FNP), fun((FNP) -> {ok, FNR} | {error, FNS})) -> {ok, list(FNR)} | {error, {list(FNR), FNS}}. 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", 552). ?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(FMS)) -> {ok, {gleam@deque:deque(FMS), 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", 576). -spec actor_builder( builder(FMX), fun((gleam@erlang@process:subject(msg(FMX))) -> FNB), integer() ) -> gleam@otp@actor:builder(state(FMX), msg(FMX), FNB). 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(), 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", 190). ?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(FLA), fun((gleam@erlang@process:subject(msg(FLA))) -> FLE), integer() ) -> gleam@otp@supervision:child_specification(FLE). 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", 180). ?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(FKX), integer()) -> gleam@otp@supervision:child_specification(gleam@erlang@process:subject(msg(FKX))). supervised(Builder, Init_timeout) -> supervised_map(Builder, fun gleam@function:identity/1, Init_timeout). -file("src/bath.gleam", 204). ?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(FLG), integer()) -> {ok, gleam@erlang@process:subject(msg(FLG))} | {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 ).