-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 => <>, 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 => <>, 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 => <>, 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 => <>, 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 => <>, 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 => <>, 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 => <>, 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 => <>, 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 ).