-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, apply/3, shutdown/3, try_map_returning/2, supervised/3, start/2]). -export_type([checkout_strategy/0, creation_strategy/0, builder/1, apply_error/0, shutdown_error/0, 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 pool(FJR) :: {pool, gleam@erlang@process:subject(msg(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_()} | {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", 63). ?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, FJZ} | {error, binary()})) -> builder(FJZ). new(Create_resource) -> {builder, 10, Create_resource, fun(_) -> nil end, f_i_f_o, lazy, false}. -file("src/bath.gleam", 77). ?DOC(" Set the pool size. Defaults to 10.\n"). -spec size(builder(FKD), integer()) -> builder(FKD). size(Builder, Size) -> _record = Builder, {builder, Size, erlang:element(3, _record), erlang:element(4, _record), erlang:element(5, _record), erlang:element(6, _record), erlang:element(7, _record)}. -file("src/bath.gleam", 85). ?DOC(" Set a shutdown function to be run for each resource when the pool exits.\n"). -spec on_shutdown(builder(FKG), fun((FKG) -> nil)) -> builder(FKG). 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", 93). ?DOC(" Change the checkout strategy for the pool. Defaults to `FIFO`.\n"). -spec checkout_strategy(builder(FKJ), checkout_strategy()) -> builder(FKJ). 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", 101). ?DOC(" Change the resource creation strategy for the pool. Defaults to `Lazy`.\n"). -spec creation_strategy(builder(FKM), creation_strategy()) -> builder(FKM). 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", 109). ?DOC(" Set whether the pool logs errors when resources fail to create.\n"). -spec log_errors(builder(FKP), boolean()) -> builder(FKP). 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", 195). ?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(FLC), gleam@erlang@process:pid_(), integer()) -> {ok, FLC} | {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", 203). -spec check_in(pool(FLG), FLG, gleam@erlang@process:pid_()) -> nil. check_in(Pool, Resource, Caller) -> gleam@erlang@process:send( erlang:element(2, Pool), {check_in, Resource, Caller} ). -file("src/bath.gleam", 218). ?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" " // Do stuff with resource...\n" " ```\n" ). -spec apply(pool(FLJ), integer(), fun((FLJ) -> FLL)) -> {ok, FLL} | {error, apply_error()}. apply(Pool, Timeout, Next) -> Self = erlang:self(), gleam@result:'try'( check_out(Pool, Self, Timeout), fun(Resource) -> Usage_result = Next(Resource), check_in(Pool, Resource, Self), {ok, Usage_result} end ). -file("src/bath.gleam", 240). ?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", 543). -spec monitor_process( gleam@erlang@process:selector(msg(FMK)), gleam@erlang@process:pid_() ) -> {gleam@erlang@process:monitor(), gleam@erlang@process:selector(msg(FMK))}. 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", 551). -spec demonitor_process( gleam@erlang@process:selector(msg(FMO)), gleam@erlang@process:monitor() ) -> gleam@erlang@process:selector(msg(FMO)). demonitor_process(Selector, Monitor) -> Selector@1 = begin _pipe = Selector, gleam@erlang@process:deselect_specific_monitor(_pipe, Monitor) end, Selector@1. -file("src/bath.gleam", 561). -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", 289). -spec handle_pool_message(state(FLS), msg(FLS)) -> gleam@otp@actor:next(state(FLS), msg(FLS)). handle_pool_message(State, Msg) -> case Msg of {check_in, Resource, Caller} -> 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 = gleam@deque:push_back( erlang:element(7, State), Resource ), gleam@otp@actor:with_selector( gleam@otp@actor:continue( 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, erlang:element(8, _record), Live_resources, Selector, erlang:element(11, _record)} end ), 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 = (erlang:element(5, State))(), gleam@result:map_error( _pipe, 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}), gleam@otp@actor:with_selector( gleam@otp@actor:continue( 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 ), Selector@1 ) end; {pool_exit, Exit_message} -> _pipe@1 = erlang:element(7, State), _pipe@2 = gleam@deque:to_list(_pipe@1), gleam@list:each(_pipe@2, erlang:element(6, State)), case erlang:element(3, Exit_message) of {abnormal, Reason} -> _pipe@3 = gleam@string:inspect(Reason), gleam@otp@actor:stop_abnormal(_pipe@3); 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@4 = erlang:element(7, State), _pipe@5 = gleam@deque:to_list(_pipe@4), gleam@list:each(_pipe@5, 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 => 418, value => _assert_fail, start => 12257, 'end' => 12329, pattern_start => 12268, pattern_end => 12314}) 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@6 = 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@7 = gleam@otp@actor:continue(_pipe@6), gleam@otp@actor:with_selector(_pipe@7, Selector@2) end end. -file("src/bath.gleam", 579). ?DOC(false). -spec try_map_returning(list(FMS), fun((FMS) -> {ok, FMU} | {error, FMV})) -> {ok, list(FMU)} | {error, {list(FMU), FMV}}. 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", 476). ?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(FLW)) -> {ok, {gleam@deque:deque(FLW), 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", 500). -spec actor_builder(builder(FMB), integer()) -> gleam@otp@actor:builder(state(FMB), msg(FMB), gleam@erlang@process:subject(msg(FMB))). 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", 161). ?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(FKS), gleam@erlang@process:subject(pool(FKS)), integer() ) -> gleam@otp@supervision:child_specification(gleam@erlang@process:subject(msg(FKS))). supervised(Builder, Pool_receiver, 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, {pool, erlang:element(3, Started)} ), {ok, Started} end ) end ). -file("src/bath.gleam", 180). ?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(FKX), integer()) -> {ok, pool(FKX)} | {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 ).