Current section
Files
Jump to
Current section
Files
src/db_pool.erl
-module(db_pool).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch, inline]).
-define(FILEPATH, "src/db_pool.gleam").
-export([new/0, size/2, on_open/2, on_close/2, on_idle/2, on_active/2, queue_target/2, queue_interval/2, start/3, supervised/3, checkout/4, checkin/3, with_connection/4, shutdown/2]).
-export_type([pool_error/1, pool/2, waiting/2, active/1, state/2, message/2]).
-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 pool_error(FQT) :: {connection_error, FQT} |
connection_timeout |
connection_unavailable.
-opaque pool(FQU, FQV) :: {pool,
integer(),
integer(),
integer(),
fun(() -> {ok, FQU} | {error, pool_error(FQV)}),
fun((FQU) -> {ok, nil} | {error, pool_error(FQV)}),
fun((FQU) -> nil),
fun((FQU) -> nil)}.
-type waiting(FQW, FQX) :: {waiting,
gleam@erlang@process:pid_(),
gleam@erlang@process:monitor(),
gleam@erlang@process:subject({ok, FQW} | {error, FQX}),
integer()}.
-type active(FQY) :: {active,
FQY,
gleam@erlang@process:monitor(),
gleam@erlang@process:timer(),
integer()}.
-type state(FQZ, FRA) :: {state,
gleam@erlang@process:subject(message(FQZ, FRA)),
integer(),
integer(),
fun(() -> {ok, FQZ} | {error, pool_error(FRA)}),
fun((FQZ) -> {ok, nil} | {error, pool_error(FRA)}),
fun((FQZ) -> nil),
fun((FQZ) -> nil),
list(FQZ),
gleam@dict:dict(gleam@erlang@process:pid_(), active(FQZ)),
rasa@queue:queue(waiting(FQZ, pool_error(FRA))),
rasa@counter:counter(),
integer(),
integer(),
integer(),
boolean(),
integer()}.
-opaque message(FRB, FRC) :: {check_out,
gleam@erlang@process:subject({ok, FRB} | {error, pool_error(FRC)}),
gleam@erlang@process:pid_(),
integer(),
integer()} |
{check_in, gleam@erlang@process:pid_(), FRB} |
{timeout, integer(), integer()} |
{deadline_expired, gleam@erlang@process:pid_(), integer()} |
{poll, integer(), integer()} |
{pool_exit, gleam@erlang@process:exit_message()} |
{caller_down, gleam@erlang@process:down()} |
{reconnect, integer()} |
{shutdown,
gleam@erlang@process:subject({ok, nil} | {error, pool_error(FRC)})}.
-file("src/db_pool.gleam", 51).
?DOC(" Returns a `Pool` that needs to be configured.\n").
-spec new() -> pool(any(), any()).
new() ->
Handle_open = fun() -> {error, connection_timeout} end,
Handle_close = fun(_) -> {ok, nil} end,
{pool,
5,
50,
1000,
Handle_open,
Handle_close,
fun(_) -> nil end,
fun(_) -> nil end}.
-file("src/db_pool.gleam", 68).
?DOC(
" Sets the size of the pool. At startup the pool will create `size`\n"
" number of connections.\n"
).
-spec size(pool(FRH, FRI), integer()) -> pool(FRH, FRI).
size(Pool, Size) ->
{pool,
Size,
erlang:element(3, Pool),
erlang:element(4, Pool),
erlang:element(5, Pool),
erlang:element(6, Pool),
erlang:element(7, Pool),
erlang:element(8, Pool)}.
-file("src/db_pool.gleam", 74).
?DOC(
" Sets the `Pool`'s `on_open` function. The provided function will be\n"
" called at startup to create connections.\n"
).
-spec on_open(pool(FRN, FRO), fun(() -> {ok, FRN} | {error, FRO})) -> pool(FRN, FRO).
on_open(Pool, Handle_open) ->
Handle_open@1 = fun() -> _pipe = Handle_open(),
gleam@result:map_error(
_pipe,
fun(Field@0) -> {connection_error, Field@0} end
) end,
{pool,
erlang:element(2, Pool),
erlang:element(3, Pool),
erlang:element(4, Pool),
Handle_open@1,
erlang:element(6, Pool),
erlang:element(7, Pool),
erlang:element(8, Pool)}.
-file("src/db_pool.gleam", 85).
?DOC(
" Sets the `Pool`'s `on_close` function. The provided function will be\n"
" called on each connection when the pool is shut down or exits.\n"
).
-spec on_close(pool(FRV, FRW), fun((FRV) -> {ok, nil} | {error, FRW})) -> pool(FRV, FRW).
on_close(Pool, Handle_close) ->
Handle_close@1 = fun(Conn) -> _pipe = Handle_close(Conn),
gleam@result:map_error(
_pipe,
fun(Field@0) -> {connection_error, Field@0} end
) end,
{pool,
erlang:element(2, Pool),
erlang:element(3, Pool),
erlang:element(4, Pool),
erlang:element(5, Pool),
Handle_close@1,
erlang:element(7, Pool),
erlang:element(8, Pool)}.
-file("src/db_pool.gleam", 101).
?DOC(
" Sets the `Pool`'s `on_idle` function. The provided function will be\n"
" called on connections when they're checked back in to the pool. If\n"
" the connection is immediately passed to a waiting caller, the callback\n"
" will not be called. The callback is also called on every connection\n"
" at startup.\n"
).
-spec on_idle(pool(FSD, FSE), fun((FSD) -> nil)) -> pool(FSD, FSE).
on_idle(Pool, Handle_idle) ->
{pool,
erlang:element(2, Pool),
erlang:element(3, Pool),
erlang:element(4, Pool),
erlang:element(5, Pool),
erlang:element(6, Pool),
Handle_idle,
erlang:element(8, Pool)}.
-file("src/db_pool.gleam", 111).
?DOC(
" Sets the `Pool`'s `on_active` function. The provided function will be\n"
" called on connections as they're removed from the pool's list of\n"
" idle connections and become active.\n"
).
-spec on_active(pool(FSJ, FSK), fun((FSJ) -> nil)) -> pool(FSJ, FSK).
on_active(Pool, Handle_active) ->
{pool,
erlang:element(2, Pool),
erlang:element(3, Pool),
erlang:element(4, Pool),
erlang:element(5, Pool),
erlang:element(6, Pool),
erlang:element(7, Pool),
Handle_active}.
-file("src/db_pool.gleam", 121).
?DOC(
" Sets the CoDel queue target in milliseconds. This is the maximum\n"
" acceptable queue delay before the pool considers itself overloaded.\n"
" Defaults to 50ms.\n"
).
-spec queue_target(pool(FSP, FSQ), integer()) -> pool(FSP, FSQ).
queue_target(Pool, Target) ->
{pool,
erlang:element(2, Pool),
Target,
erlang:element(4, Pool),
erlang:element(5, Pool),
erlang:element(6, Pool),
erlang:element(7, Pool),
erlang:element(8, Pool)}.
-file("src/db_pool.gleam", 128).
?DOC(
" Sets the CoDel queue interval in milliseconds. This is the length\n"
" of each CoDel measurement interval. The pool evaluates queue health\n"
" at each interval boundary. Defaults to 1000ms.\n"
).
-spec queue_interval(pool(FSV, FSW), integer()) -> pool(FSV, FSW).
queue_interval(Pool, Interval) ->
{pool,
erlang:element(2, Pool),
erlang:element(3, Pool),
Interval,
erlang:element(5, Pool),
erlang:element(6, Pool),
erlang:element(7, Pool),
erlang:element(8, Pool)}.
-file("src/db_pool.gleam", 940).
-spec close_idle(state(any(), any())) -> nil.
close_idle(State) ->
gleam@list:each(
erlang:element(9, State),
fun(Conn) ->
_ = (erlang:element(6, State))(Conn),
nil
end
).
-file("src/db_pool.gleam", 931).
-spec close_active(state(any(), any())) -> nil.
close_active(State) ->
gleam@dict:each(
erlang:element(10, State),
fun(_, Active) ->
_ = gleam@erlang@process:cancel_timer(erlang:element(4, Active)),
gleam@erlang@process:demonitor_process(erlang:element(3, Active)),
_ = (erlang:element(6, State))(erlang:element(2, Active)),
nil
end
).
-file("src/db_pool.gleam", 916).
-spec drop_waiter(waiting(any(), pool_error(any()))) -> nil.
drop_waiter(Waiting) ->
gleam@otp@actor:send(
erlang:element(4, Waiting),
{error, connection_unavailable}
),
gleam@erlang@process:demonitor_process(erlang:element(3, Waiting)).
-file("src/db_pool.gleam", 921).
-spec drain_queue(state(any(), any())) -> nil.
drain_queue(State) ->
case rasa@queue:pop(erlang:element(11, State)) of
{ok, Waiting} ->
drop_waiter(Waiting),
drain_queue(State);
_ ->
nil
end.
-file("src/db_pool.gleam", 693).
-spec schedule_reconnect(state(any(), any()), integer()) -> nil.
schedule_reconnect(State, Backoff) ->
Half = Backoff div 2,
Delay = Half + gleam@int:random(Half + 1),
Next_backoff = gleam@int:min(Backoff * 2, 30000),
_ = gleam@erlang@process:send_after(
erlang:element(2, State),
Delay,
{reconnect, Next_backoff}
),
nil.
-file("src/db_pool.gleam", 803).
-spec serve_waiter(state(FZK, FZL), waiting(FZK, pool_error(FZL)), FZK) -> state(FZK, FZL).
serve_waiter(State, Waiting, Conn) ->
case erlang:is_process_alive(erlang:element(2, Waiting)) of
false ->
gleam@erlang@process:demonitor_process(erlang:element(3, Waiting)),
Now = rasa@counter:next(erlang:element(12, State)),
codel_dequeue(State, Now, Conn);
true ->
Now@1 = rasa@counter:next(erlang:element(12, State)),
Deadline_timer = gleam@erlang@process:send_after(
erlang:element(2, State),
erlang:element(5, Waiting),
{deadline_expired, erlang:element(2, Waiting), Now@1}
),
Activated = {active,
Conn,
erlang:element(3, Waiting),
Deadline_timer,
Now@1},
Active = gleam@dict:insert(
erlang:element(10, State),
erlang:element(2, Waiting),
Activated
),
gleam@erlang@process:send(erlang:element(4, Waiting), {ok, Conn}),
(erlang:element(8, State))(Conn),
{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),
Active,
erlang:element(11, State),
erlang:element(12, State),
erlang:element(13, State),
erlang:element(14, State),
erlang:element(15, State),
erlang:element(16, State),
erlang:element(17, State)}
end.
-file("src/db_pool.gleam", 769).
-spec dequeue_slow(state(FZE, FZF), integer(), integer(), FZE) -> state(FZE, FZF).
dequeue_slow(State, Now, Timeout, Conn) ->
case rasa@queue:first(erlang:element(11, State)) of
{ok, {Sent, Waiting}} when (Now - Sent) > Timeout ->
case rasa@queue:delete(erlang:element(11, State), Sent) of
{ok, nil} -> nil;
_assert_fail ->
erlang:error(#{gleam_error => let_assert,
message => <<"Pattern match failed, no pattern matched the value."/utf8>>,
file => <<?FILEPATH/utf8>>,
module => <<"db_pool"/utf8>>,
function => <<"dequeue_slow"/utf8>>,
line => 777,
value => _assert_fail,
start => 22614,
'end' => 22666,
pattern_start => 22625,
pattern_end => 22632})
end,
drop_waiter(Waiting),
_pipe = State,
dequeue_slow(_pipe, Now, Timeout, Conn);
{ok, {Sent@1, Waiting@1}} ->
case rasa@queue:delete(erlang:element(11, State), Sent@1) of
{ok, nil} -> nil;
_assert_fail@1 ->
erlang:error(#{gleam_error => let_assert,
message => <<"Pattern match failed, no pattern matched the value."/utf8>>,
file => <<?FILEPATH/utf8>>,
module => <<"db_pool"/utf8>>,
function => <<"dequeue_slow"/utf8>>,
line => 785,
value => _assert_fail@1,
start => 22792,
'end' => 22844,
pattern_start => 22803,
pattern_end => 22810})
end,
Delay = Now - Sent@1,
State@1 = case Delay < erlang:element(15, State) of
true ->
{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),
erlang:element(10, State),
erlang:element(11, State),
erlang:element(12, State),
erlang:element(13, State),
erlang:element(14, State),
Delay,
erlang:element(16, State),
erlang:element(17, State)};
false ->
State
end,
serve_waiter(State@1, Waiting@1, Conn);
_ ->
(erlang:element(7, State))(Conn),
{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),
[Conn | erlang:element(9, State)],
erlang:element(10, State),
erlang:element(11, State),
erlang:element(12, State),
erlang:element(13, State),
erlang:element(14, State),
erlang:element(15, State),
erlang:element(16, State),
erlang:element(17, State)}
end.
-file("src/db_pool.gleam", 744).
-spec dequeue_fast(state(FYY, FYZ), integer(), FYY) -> state(FYY, FYZ).
dequeue_fast(State, Now, Conn) ->
case rasa@queue:first(erlang:element(11, State)) of
{ok, {Sent, Waiting}} ->
case rasa@queue:delete(erlang:element(11, State), Sent) of
{ok, nil} -> nil;
_assert_fail ->
erlang:error(#{gleam_error => let_assert,
message => <<"Pattern match failed, no pattern matched the value."/utf8>>,
file => <<?FILEPATH/utf8>>,
module => <<"db_pool"/utf8>>,
function => <<"dequeue_fast"/utf8>>,
line => 751,
value => _assert_fail,
start => 21937,
'end' => 21989,
pattern_start => 21948,
pattern_end => 21955})
end,
Delay = Now - Sent,
State@1 = case Delay < erlang:element(15, State) of
true ->
{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),
erlang:element(10, State),
erlang:element(11, State),
erlang:element(12, State),
erlang:element(13, State),
erlang:element(14, State),
Delay,
erlang:element(16, State),
erlang:element(17, State)};
false ->
State
end,
serve_waiter(State@1, Waiting, Conn);
_ ->
(erlang:element(7, State))(Conn),
{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),
[Conn | erlang:element(9, State)],
erlang:element(10, State),
erlang:element(11, State),
erlang:element(12, State),
erlang:element(13, State),
erlang:element(14, State),
erlang:element(15, State),
erlang:element(16, State),
erlang:element(17, State)}
end.
-file("src/db_pool.gleam", 718).
-spec dequeue_first(state(FYS, FYT), integer(), FYS) -> state(FYS, FYT).
dequeue_first(State, Now, Conn) ->
Next = Now + erlang:element(14, State),
Slow = erlang:element(15, State) > erlang:element(13, State),
case rasa@queue:first(erlang:element(11, State)) of
{ok, {Sent, Waiting}} ->
case rasa@queue:delete(erlang:element(11, State), Sent) of
{ok, nil} -> nil;
_assert_fail ->
erlang:error(#{gleam_error => let_assert,
message => <<"Pattern match failed, no pattern matched the value."/utf8>>,
file => <<?FILEPATH/utf8>>,
module => <<"db_pool"/utf8>>,
function => <<"dequeue_first"/utf8>>,
line => 728,
value => _assert_fail,
start => 21396,
'end' => 21448,
pattern_start => 21407,
pattern_end => 21414})
end,
Delay = Now - Sent,
State@1 = {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),
erlang:element(10, State),
erlang:element(11, State),
erlang:element(12, State),
erlang:element(13, State),
erlang:element(14, State),
Delay,
Slow,
Next},
serve_waiter(State@1, Waiting, Conn);
_ ->
(erlang:element(7, State))(Conn),
{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),
[Conn | erlang:element(9, State)],
erlang:element(10, State),
erlang:element(11, State),
erlang:element(12, State),
erlang:element(13, State),
erlang:element(14, State),
0,
Slow,
Next}
end.
-file("src/db_pool.gleam", 705).
-spec codel_dequeue(state(FYM, FYN), integer(), FYM) -> state(FYM, FYN).
codel_dequeue(State, Now, Conn) ->
case {(Now >= erlang:element(17, State)), erlang:element(16, State)} of
{true, _} ->
dequeue_first(State, Now, Conn);
{false, false} ->
dequeue_fast(State, Now, Conn);
{false, true} ->
dequeue_slow(State, Now, erlang:element(13, State) * 2, Conn)
end.
-file("src/db_pool.gleam", 679).
?DOC(
" Called when a reconnect timer fires. Attempts to open a replacement\n"
" connection. On success, the connection is fed through CoDel to serve\n"
" a waiter or return to idle. On failure, another reconnect is scheduled\n"
" with increased backoff (randomized exponential, capped at 30s).\n"
).
-spec do_reconnect(state(FYC, FYD), integer()) -> state(FYC, FYD).
do_reconnect(State, Backoff) ->
case (erlang:element(5, State))() of
{ok, Conn} ->
State@1 = {state,
erlang:element(2, State),
erlang:element(3, State),
erlang:element(4, State) + 1,
erlang:element(5, State),
erlang:element(6, State),
erlang:element(7, State),
erlang:element(8, State),
erlang:element(9, State),
erlang:element(10, State),
erlang:element(11, State),
erlang:element(12, State),
erlang:element(13, State),
erlang:element(14, State),
erlang:element(15, State),
erlang:element(16, State),
erlang:element(17, State)},
Now = rasa@counter:next(erlang:element(12, State@1)),
codel_dequeue(State@1, Now, Conn);
_ ->
schedule_reconnect(State, Backoff),
State
end.
-file("src/db_pool.gleam", 898).
-spec start_poll(state(GAL, GAM), integer(), integer()) -> state(GAL, GAM).
start_poll(State, Now, Last_sent) ->
Poll_time = Now + erlang:element(14, State),
_ = gleam@erlang@process:send_after(
erlang:element(2, State),
case 1000000 of
0 -> 0;
Gleam@denominator -> erlang:element(14, State) div Gleam@denominator
end,
{poll, Poll_time, Last_sent}
),
State.
-file("src/db_pool.gleam", 880).
-spec poll_drop_slow(state(GAF, GAG), integer(), integer()) -> state(GAF, GAG).
poll_drop_slow(State, Now, Timeout) ->
case rasa@queue:first(erlang:element(11, State)) of
{ok, {Sent, Waiting}} when (Now - Sent) > Timeout ->
case rasa@queue:delete(erlang:element(11, State), Sent) of
{ok, nil} -> nil;
_assert_fail ->
erlang:error(#{gleam_error => let_assert,
message => <<"Pattern match failed, no pattern matched the value."/utf8>>,
file => <<?FILEPATH/utf8>>,
module => <<"db_pool"/utf8>>,
function => <<"poll_drop_slow"/utf8>>,
line => 887,
value => _assert_fail,
start => 25079,
'end' => 25131,
pattern_start => 25090,
pattern_end => 25097})
end,
drop_waiter(Waiting),
_pipe = State,
poll_drop_slow(_pipe, Now, Timeout);
_ ->
State
end.
-file("src/db_pool.gleam", 864).
-spec codel_timeout(state(FZZ, GAA), integer(), integer()) -> state(FZZ, GAA).
codel_timeout(State, Delay, Time) ->
case {Time >= erlang:element(17, State),
erlang:element(15, State) > erlang:element(13, State)} of
{true, true} ->
_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),
erlang:element(9, State),
erlang:element(10, State),
erlang:element(11, State),
erlang:element(12, State),
erlang:element(13, State),
erlang:element(14, State),
Delay,
true,
Time + erlang:element(14, State)},
poll_drop_slow(_pipe, Time, erlang:element(13, State) * 2);
{true, false} ->
{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),
erlang:element(10, State),
erlang:element(11, State),
erlang:element(12, State),
erlang:element(13, State),
erlang:element(14, State),
Delay,
false,
Time + erlang:element(14, State)};
{_, _} ->
State
end.
-file("src/db_pool.gleam", 846).
-spec do_poll(state(FZT, FZU), integer(), integer()) -> state(FZT, FZU).
do_poll(State, Time, Last_sent) ->
case rasa@queue:first(erlang:element(11, State)) of
{ok, {Sent, _}} when Sent =< Last_sent ->
Delay = Time - Sent,
_pipe = State,
_pipe@1 = codel_timeout(_pipe, Delay, Time),
start_poll(_pipe@1, Time, Sent);
{ok, {Sent@1, _}} ->
start_poll(State, Time, Sent@1);
_ ->
start_poll(State, Time, Time)
end.
-file("src/db_pool.gleam", 607).
?DOC(
" Called when a caller process dies while holding a connection or waiting.\n"
" If the caller held an active connection, the connection is closed and\n"
" replaced. If the caller was waiting in the queue, the entry is cleaned\n"
" up lazily: `serve_waiter` checks `process.is_alive` at dequeue time,\n"
" and `do_expire` removes entries when their timeout fires. The queue is\n"
" keyed by timestamp, so there is no efficient PID-based removal.\n"
).
-spec do_caller_down(state(FXQ, FXR), gleam@erlang@process:pid_()) -> state(FXQ, FXR).
do_caller_down(State, Pid) ->
case gleam_stdlib:map_get(erlang:element(10, State), Pid) of
{ok, Prev} ->
_ = gleam@erlang@process:cancel_timer(erlang:element(4, Prev)),
gleam@erlang@process:demonitor_process(erlang:element(3, Prev)),
Active = gleam@dict:delete(erlang:element(10, State), Pid),
State@1 = {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),
Active,
erlang:element(11, State),
erlang:element(12, State),
erlang:element(13, State),
erlang:element(14, State),
erlang:element(15, State),
erlang:element(16, State),
erlang:element(17, State)},
_ = (erlang:element(6, State@1))(erlang:element(2, Prev)),
State@2 = {state,
erlang:element(2, State@1),
erlang:element(3, State@1),
erlang:element(4, State@1) - 1,
erlang:element(5, State@1),
erlang:element(6, State@1),
erlang:element(7, State@1),
erlang:element(8, State@1),
erlang:element(9, State@1),
erlang:element(10, State@1),
erlang:element(11, State@1),
erlang:element(12, State@1),
erlang:element(13, State@1),
erlang:element(14, State@1),
erlang:element(15, State@1),
erlang:element(16, State@1),
erlang:element(17, State@1)},
case (erlang:element(5, State@2))() of
{ok, Conn} ->
State@3 = {state,
erlang:element(2, State@2),
erlang:element(3, State@2),
erlang:element(4, State@2) + 1,
erlang:element(5, State@2),
erlang:element(6, State@2),
erlang:element(7, State@2),
erlang:element(8, State@2),
erlang:element(9, State@2),
erlang:element(10, State@2),
erlang:element(11, State@2),
erlang:element(12, State@2),
erlang:element(13, State@2),
erlang:element(14, State@2),
erlang:element(15, State@2),
erlang:element(16, State@2),
erlang:element(17, State@2)},
Now = rasa@counter:next(erlang:element(12, State@3)),
codel_dequeue(State@3, Now, Conn);
_ ->
schedule_reconnect(State@2, 1000),
State@2
end;
_ ->
State
end.
-file("src/db_pool.gleam", 640).
-spec do_deadline_expired(
state(FXW, FXX),
gleam@erlang@process:pid_(),
integer()
) -> state(FXW, FXX).
do_deadline_expired(State, Caller, Checkout_time) ->
_pipe = gleam_stdlib:map_get(erlang:element(10, State), Caller),
_pipe@1 = gleam@result:map(
_pipe,
fun(Active) ->
gleam@bool:guard(
erlang:element(5, Active) /= Checkout_time,
State,
fun() ->
gleam@erlang@process:demonitor_process(
erlang:element(3, Active)
),
Active_dict = gleam@dict:delete(
erlang:element(10, State),
Caller
),
State@1 = {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),
Active_dict,
erlang:element(11, State),
erlang:element(12, State),
erlang:element(13, State),
erlang:element(14, State),
erlang:element(15, State),
erlang:element(16, State),
erlang:element(17, State)},
_ = (erlang:element(6, State@1))(erlang:element(2, Active)),
State@2 = {state,
erlang:element(2, State@1),
erlang:element(3, State@1),
erlang:element(4, State@1) - 1,
erlang:element(5, State@1),
erlang:element(6, State@1),
erlang:element(7, State@1),
erlang:element(8, State@1),
erlang:element(9, State@1),
erlang:element(10, State@1),
erlang:element(11, State@1),
erlang:element(12, State@1),
erlang:element(13, State@1),
erlang:element(14, State@1),
erlang:element(15, State@1),
erlang:element(16, State@1),
erlang:element(17, State@1)},
case (erlang:element(5, State@2))() of
{ok, Conn} ->
State@3 = {state,
erlang:element(2, State@2),
erlang:element(3, State@2),
erlang:element(4, State@2) + 1,
erlang:element(5, State@2),
erlang:element(6, State@2),
erlang:element(7, State@2),
erlang:element(8, State@2),
erlang:element(9, State@2),
erlang:element(10, State@2),
erlang:element(11, State@2),
erlang:element(12, State@2),
erlang:element(13, State@2),
erlang:element(14, State@2),
erlang:element(15, State@2),
erlang:element(16, State@2),
erlang:element(17, State@2)},
Now = rasa@counter:next(erlang:element(12, State@3)),
codel_dequeue(State@3, Now, Conn);
_ ->
schedule_reconnect(State@2, 1000),
State@2
end
end
)
end
),
gleam@result:unwrap(_pipe@1, State).
-file("src/db_pool.gleam", 568).
-spec do_expire(state(FXK, FXL), integer(), integer()) -> state(FXK, FXL).
do_expire(State, Sent, Timeout) ->
_pipe = rasa@queue:at(erlang:element(11, State), Sent),
_pipe@1 = gleam@result:map(
_pipe,
fun(Waiting) ->
Now = rasa@counter:next(erlang:element(12, State)),
gleam@bool:lazy_guard(
(Now < (Sent + (Timeout * 1000000))),
fun() ->
Remaining_ns = (Sent + (Timeout * 1000000)) - Now,
Remaining_ms = case 1000000 of
0 -> 0;
Gleam@denominator -> Remaining_ns div Gleam@denominator
end,
_ = gleam@erlang@process:send_after(
erlang:element(2, State),
Remaining_ms,
{timeout, Sent, Timeout}
),
State
end,
fun() ->
case rasa@queue:delete(erlang:element(11, State), Sent) of
{ok, nil} -> nil;
_assert_fail ->
erlang:error(#{gleam_error => let_assert,
message => <<"Pattern match failed, no pattern matched the value."/utf8>>,
file => <<?FILEPATH/utf8>>,
module => <<"db_pool"/utf8>>,
function => <<"do_expire"/utf8>>,
line => 590,
value => _assert_fail,
start => 16973,
'end' => 17025,
pattern_start => 16984,
pattern_end => 16991})
end,
gleam@otp@actor:send(
erlang:element(4, Waiting),
{error, connection_timeout}
),
gleam@erlang@process:demonitor_process(
erlang:element(3, Waiting)
),
State
end
)
end
),
gleam@result:unwrap(_pipe@1, State).
-file("src/db_pool.gleam", 548).
-spec do_enqueue(
state(FXA, FXB),
gleam@erlang@process:pid_(),
gleam@erlang@process:subject({ok, FXA} | {error, pool_error(FXB)}),
integer(),
integer()
) -> state(FXA, FXB).
do_enqueue(State, Caller, Client, Timeout, Deadline) ->
Monitor = gleam@erlang@process:monitor(Caller),
Waiting = {waiting, Caller, Monitor, Client, Deadline},
Sent_at@1 = case rasa@queue:push(erlang:element(11, State), Waiting) of
{ok, Sent_at} -> Sent_at;
_assert_fail ->
erlang:error(#{gleam_error => let_assert,
message => <<"Pattern match failed, no pattern matched the value."/utf8>>,
file => <<?FILEPATH/utf8>>,
module => <<"db_pool"/utf8>>,
function => <<"do_enqueue"/utf8>>,
line => 560,
value => _assert_fail,
start => 16268,
'end' => 16325,
pattern_start => 16279,
pattern_end => 16290})
end,
_ = gleam@erlang@process:send_after(
erlang:element(2, State),
Timeout,
{timeout, Sent_at@1, Timeout}
),
State.
-file("src/db_pool.gleam", 473).
?DOC(
" Try to check out a connection. Returns Ok(state) if served\n"
" (re-entrant checkout or idle conn available), Error(Nil) if the\n"
" caller should be enqueued.\n"
).
-spec do_checkout(
state(FWI, FWJ),
gleam@erlang@process:pid_(),
gleam@erlang@process:subject({ok, FWI} | {error, pool_error(FWJ)}),
integer()
) -> {ok, state(FWI, FWJ)} | {error, nil}.
do_checkout(State, Caller, Client, Deadline) ->
case gleam_stdlib:map_get(erlang:element(10, State), Caller) of
{ok, Active} ->
gleam@otp@actor:send(Client, {ok, erlang:element(2, Active)}),
{ok, State};
_ ->
case erlang:element(9, State) of
[Conn | Rest] ->
Monitor = gleam@erlang@process:monitor(Caller),
Now = rasa@counter:next(erlang:element(12, State)),
Deadline_timer = gleam@erlang@process:send_after(
erlang:element(2, State),
Deadline,
{deadline_expired, Caller, Now}
),
Activated = {active, Conn, Monitor, Deadline_timer, Now},
Active@1 = gleam@dict:insert(
erlang:element(10, State),
Caller,
Activated
),
(erlang:element(8, State))(Conn),
gleam@otp@actor:send(Client, {ok, Conn}),
{ok,
{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),
Rest,
Active@1,
erlang:element(11, State),
erlang:element(12, State),
erlang:element(13, State),
erlang:element(14, State),
erlang:element(15, State),
erlang:element(16, State),
erlang:element(17, State)}};
[] ->
{error, nil}
end
end.
-file("src/db_pool.gleam", 521).
?DOC(
" Called when a client returns a connection to the pool.\n"
" Cleans up monitoring/deadline, then either serves a waiter\n"
" via CoDel or returns the connection to idle.\n"
).
-spec do_checkin(state(FWU, FWV), gleam@erlang@process:pid_(), FWU) -> state(FWU, FWV).
do_checkin(State, Caller, Conn) ->
case gleam_stdlib:map_get(erlang:element(10, State), Caller) of
{ok, Prev} ->
case erlang:element(2, Prev) =:= Conn of
true ->
nil;
false ->
_pipe = <<"(db_pool) unexpected connection checked in for the current process"/utf8>>,
gleam_stdlib:println_error(_pipe)
end,
_ = gleam@erlang@process:cancel_timer(erlang:element(4, Prev)),
gleam@erlang@process:demonitor_process(erlang:element(3, Prev)),
Active = gleam@dict:delete(erlang:element(10, State), Caller),
State@1 = {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),
Active,
erlang:element(11, State),
erlang:element(12, State),
erlang:element(13, State),
erlang:element(14, State),
erlang:element(15, State),
erlang:element(16, State),
erlang:element(17, State)},
Now = rasa@counter:next(erlang:element(12, State@1)),
codel_dequeue(State@1, Now, erlang:element(2, Prev));
_ ->
State
end.
-file("src/db_pool.gleam", 402).
-spec handle_message(state(FVW, FVX), message(FVW, FVX)) -> gleam@otp@actor:next(state(FVW, FVX), message(FVW, FVX)).
handle_message(State, Msg) ->
case Msg of
{check_in, Caller, Conn} ->
State@1 = do_checkin(State, Caller, Conn),
gleam@otp@actor:continue(State@1);
{check_out, Client, Caller@1, Timeout, Deadline} ->
State@2 = begin
_pipe = do_checkout(State, Caller@1, Client, Deadline),
gleam@result:lazy_unwrap(
_pipe,
fun() ->
do_enqueue(State, Caller@1, Client, Timeout, Deadline)
end
)
end,
gleam@otp@actor:continue(State@2);
{timeout, Time_sent, Timeout@1} ->
State@3 = do_expire(State, Time_sent, Timeout@1),
gleam@otp@actor:continue(State@3);
{deadline_expired, Caller@2, Checkout_time} ->
_pipe@1 = State,
_pipe@2 = do_deadline_expired(_pipe@1, Caller@2, Checkout_time),
gleam@otp@actor:continue(_pipe@2);
{caller_down, Down} ->
Pid@1 = case Down of
{process_down, _, Pid, _} -> Pid;
_assert_fail ->
erlang:error(#{gleam_error => let_assert,
message => <<"Pattern match failed, no pattern matched the value."/utf8>>,
file => <<?FILEPATH/utf8>>,
module => <<"db_pool"/utf8>>,
function => <<"handle_message"/utf8>>,
line => 430,
value => _assert_fail,
start => 12587,
'end' => 12634,
pattern_start => 12598,
pattern_end => 12627})
end,
_pipe@3 = State,
_pipe@4 = do_caller_down(_pipe@3, Pid@1),
gleam@otp@actor:continue(_pipe@4);
{poll, Time, Last_sent} ->
_pipe@5 = State,
_pipe@6 = do_poll(_pipe@5, Time, Last_sent),
gleam@otp@actor:continue(_pipe@6);
{reconnect, Backoff} ->
_pipe@7 = State,
_pipe@8 = do_reconnect(_pipe@7, Backoff),
gleam@otp@actor:continue(_pipe@8);
{pool_exit, Exit} ->
drain_queue(State),
_ = close_active(State),
close_idle(State),
case erlang:element(3, Exit) of
normal ->
gleam@otp@actor:stop();
killed ->
gleam@otp@actor:stop_abnormal(<<"pool killed"/utf8>>);
{abnormal, _} ->
gleam@otp@actor:stop_abnormal(
<<"pool stopped abnormally"/utf8>>
)
end;
{shutdown, Client@1} ->
drain_queue(State),
_ = close_active(State),
close_idle(State),
gleam@otp@actor:send(Client@1, {ok, nil}),
gleam@otp@actor:stop()
end.
-file("src/db_pool.gleam", 333).
-spec initialise_pool(
gleam@erlang@process:subject(message(FVD, FVE)),
pool(FVD, FVE),
rasa@counter:counter()
) -> {ok,
gleam@otp@actor:initialised(state(FVD, FVE), message(FVD, FVE), gleam@erlang@process:subject(message(FVD, FVE)))} |
{error, binary()}.
initialise_pool(Self, Pool, Counter) ->
gleam_erlang_ffi:trap_exits(true),
Selector = begin
_pipe = gleam_erlang_ffi:new_selector(),
_pipe@1 = gleam@erlang@process:select(_pipe, Self),
_pipe@2 = gleam@erlang@process:select_trapped_exits(
_pipe@1,
fun(Field@0) -> {pool_exit, Field@0} end
),
gleam@erlang@process:select_monitors(
_pipe@2,
fun(Field@0) -> {caller_down, Field@0} end
)
end,
Connections = begin
_pipe@3 = gleam@list:repeat(<<""/utf8>>, erlang:element(2, Pool)),
_pipe@4 = gleam@list:try_map(
_pipe@3,
fun(_) -> (erlang:element(5, Pool))() end
),
gleam@result:map_error(
_pipe@4,
fun(_) -> <<"(db_pool) Failed to open connections"/utf8>> end
)
end,
gleam@result:map(
Connections,
fun(Conns) ->
gleam@list:each(Conns, erlang:element(7, Pool)),
Q = begin
_pipe@5 = rasa@queue:new(),
_pipe@6 = rasa@queue:with_access(_pipe@5, private),
_pipe@7 = rasa@queue:with_counter(_pipe@6, Counter),
rasa@queue:build(_pipe@7)
end,
Now = rasa@counter:next(Counter),
_ = gleam@erlang@process:send_after(
Self,
erlang:element(4, Pool),
{poll, Now, Now}
),
State = {state,
Self,
erlang:element(2, Pool),
erlang:element(2, Pool),
erlang:element(5, Pool),
erlang:element(6, Pool),
erlang:element(7, Pool),
erlang:element(8, Pool),
Conns,
maps:new(),
Q,
Counter,
erlang:element(3, Pool) * 1000000,
erlang:element(4, Pool) * 1000000,
0,
false,
Now + (erlang:element(4, Pool) * 1000000)},
_pipe@8 = gleam@otp@actor:initialised(State),
_pipe@9 = gleam@otp@actor:selecting(_pipe@8, Selector),
gleam@otp@actor:returning(_pipe@9, Self)
end
).
-file("src/db_pool.gleam", 181).
?DOC(
" Starts a connection pool and registers it under `name`. All\n"
" configured connections are opened eagerly during initialisation.\n"
"\n"
" The `timeout` parameter is the maximum time in milliseconds allowed\n"
" for the actor to initialise (open all connections).\n"
"\n"
" The pool actor traps exits so it can perform cleanup when its\n"
" parent or linked processes terminate.\n"
).
-spec start(
pool(FTB, FTC),
gleam@erlang@process:name(message(FTB, FTC)),
integer()
) -> {ok,
gleam@otp@actor:started(gleam@erlang@process:subject(message(FTB, FTC)))} |
{error, gleam@otp@actor:start_error()}.
start(Pool, Name, Timeout) ->
Counter = rasa@counter:monotonic_time(nanosecond),
_pipe = gleam@otp@actor:new_with_initialiser(
Timeout,
fun(_capture) -> initialise_pool(_capture, Pool, Counter) end
),
_pipe@1 = gleam@otp@actor:on_message(_pipe, fun handle_message/2),
_pipe@2 = gleam@otp@actor:named(_pipe@1, Name),
gleam@otp@actor:start(_pipe@2).
-file("src/db_pool.gleam", 201).
?DOC(
" Creates a `supervision.ChildSpecification` so the pool can be\n"
" added to an application's supervision tree.\n"
"\n"
" The `timeout` parameter is used for both the actor initialisation\n"
" timeout and the supervisor's shutdown timeout. The restart strategy\n"
" is set to `Transient` — the pool is restarted only if it terminates\n"
" abnormally.\n"
).
-spec supervised(
pool(FTO, FTP),
gleam@erlang@process:name(message(FTO, FTP)),
integer()
) -> gleam@otp@supervision:child_specification(gleam@erlang@process:subject(message(FTO, FTP))).
supervised(Pool, Name, Timeout) ->
_pipe = gleam@otp@supervision:worker(
fun() -> start(Pool, Name, Timeout) end
),
_pipe@1 = gleam@otp@supervision:timeout(_pipe, Timeout),
gleam@otp@supervision:restart(_pipe@1, transient).
-file("src/db_pool.gleam", 254).
?DOC(
" Checks out a connection from the pool.\n"
"\n"
" The `caller` should be `process.self()` of the calling process. The\n"
" pool monitors this process and reclaims the connection if it crashes.\n"
"\n"
" If a connection is available it is returned immediately. If all\n"
" connections are in use the caller is added to a FIFO queue and will\n"
" receive a connection when one becomes available, or a\n"
" `ConnectionTimeout` error after `timeout` milliseconds.\n"
"\n"
" The `deadline` parameter sets the maximum time in milliseconds that\n"
" the connection may be held. If the caller has not checked in by then,\n"
" the pool forcibly closes the connection, replaces it, and the caller\n"
" is left holding a now-closed connection.\n"
"\n"
" Re-entrant: calling `checkout` again from the same process returns\n"
" the already checked-out connection. The original deadline is\n"
" preserved — a second checkout cannot extend it.\n"
"\n"
" Panics if the pool actor is unreachable.\n"
).
-spec checkout(
gleam@erlang@process:subject(message(FTZ, FUA)),
gleam@erlang@process:pid_(),
integer(),
integer()
) -> {ok, FTZ} | {error, pool_error(FUA)}.
checkout(Pool, Caller, Timeout, Deadline) ->
gleam@erlang@process:call(
Pool,
Timeout + 5000,
fun(_capture) -> {check_out, _capture, Caller, Timeout, Deadline} end
).
-file("src/db_pool.gleam", 274).
?DOC(
" Returns a connection back to the pool.\n"
"\n"
" Expects the `conn` value to be the same connection that was originally\n"
" checked out, and `caller` should be the `Pid` that checked it out.\n"
" If the caller has no active connection the checkin is silently\n"
" ignored.\n"
).
-spec checkin(
gleam@erlang@process:subject(message(FUH, any())),
FUH,
gleam@erlang@process:pid_()
) -> nil.
checkin(Pool, Conn, Caller) ->
gleam@erlang@process:send(Pool, {check_in, Caller, Conn}).
-file("src/db_pool.gleam", 291).
?DOC(
" Checks out a connection from the pool and passes it to the provided\n"
" callback function. The connection is automatically checked back in\n"
" after the callback function returns.\n"
"\n"
" If the callback panics, the connection is not checked in immediately.\n"
" It is reclaimed when the caller process exits (via the pool's\n"
" monitor).\n"
"\n"
" Panics if the pool actor is unreachable (crashed or shut down).\n"
).
-spec with_connection(
gleam@erlang@process:subject(message(FUM, FUN)),
integer(),
integer(),
fun((FUM) -> FUR)
) -> {ok, FUR} | {error, pool_error(FUN)}.
with_connection(Pool, Timeout, Deadline, Next) ->
Caller = erlang:self(),
_pipe = gleam@erlang@process:call(
Pool,
Timeout + 5000,
fun(_capture) -> {check_out, _capture, Caller, Timeout, Deadline} end
),
gleam@result:map(
_pipe,
fun(Conn) ->
Res = Next(Conn),
gleam@erlang@process:send(Pool, {check_in, Caller, Conn}),
Res
end
).
-file("src/db_pool.gleam", 324).
?DOC(
" Shuts down the pool gracefully within `timeout` milliseconds.\n"
"\n"
" All waiting callers in the queue are drained and sent a\n"
" `ConnectionUnavailable` error. Active (checked-out) connections\n"
" are closed, their deadline timers cancelled, and their monitors\n"
" removed. Idle connections are then closed via the configured\n"
" `on_close` callback.\n"
"\n"
" Panics if the pool actor is unreachable or does not respond\n"
" within the timeout.\n"
).
-spec shutdown(gleam@erlang@process:subject(message(any(), FUW)), integer()) -> {ok,
nil} |
{error, pool_error(FUW)}.
shutdown(Pool, Timeout) ->
gleam@erlang@process:call(
Pool,
Timeout + 5000,
fun(Field@0) -> {shutdown, Field@0} end
).