Current section
Files
Jump to
Current section
Files
src/postgleam@pool.erl
-module(postgleam@pool).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch, inline]).
-define(FILEPATH, "src/postgleam/pool.gleam").
-export([start/2, 'query'/4, simple_query/3, shutdown/2, query_with/5]).
-export_type([pool_message/0, pool_state/0, pool_response/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 pool_message() :: {execute,
fun((postgleam@connection:connection_state(), gleam@dict:dict(integer(), postgleam@codec:codec()), postgleam@config:config()) -> {{ok,
postgleam@connection:extended_query_result()} |
{error, postgleam@error:error()},
postgleam@connection:connection_state()}),
gleam@erlang@process:subject({ok,
postgleam@connection:extended_query_result()} |
{error, postgleam@error:error()})} |
{simple_execute,
fun((postgleam@connection:connection_state(), postgleam@config:config()) -> {{ok,
list(postgleam@connection:simple_query_result())} |
{error, postgleam@error:error()},
postgleam@connection:connection_state()}),
gleam@erlang@process:subject({ok,
list(postgleam@connection:simple_query_result())} |
{error, postgleam@error:error()})} |
{shutdown, gleam@erlang@process:subject(nil)}.
-type pool_state() :: {pool_state,
list(postgleam@connection:connection_state()),
postgleam@config:config(),
gleam@dict:dict(integer(), postgleam@codec:codec()),
integer()}.
-type pool_response(JFD) :: {pool_response, list(JFD), integer(), binary()}.
-file("src/postgleam/pool.gleam", 103).
-spec disconnect_all(list(postgleam@connection:connection_state())) -> nil.
disconnect_all(Conns) ->
case Conns of
[] ->
nil;
[Conn | Rest] ->
postgleam@connection:disconnect(Conn),
disconnect_all(Rest)
end.
-file("src/postgleam/pool.gleam", 84).
-spec connect_pool(
postgleam@config:config(),
integer(),
list(postgleam@connection:connection_state())
) -> {ok, list(postgleam@connection:connection_state())} |
{error, postgleam@error:error()}.
connect_pool(Config, Remaining, Acc) ->
case Remaining of
0 ->
{ok, Acc};
_ ->
case postgleam@connection:connect(Config) of
{ok, Conn} ->
connect_pool(Config, Remaining - 1, [Conn | Acc]);
{error, E} ->
disconnect_all(Acc),
{error, E}
end
end.
-file("src/postgleam/pool.gleam", 157).
-spec append(list(JFP), list(JFP)) -> list(JFP).
append(A, B) ->
case A of
[] ->
B;
[X | Rest] ->
[X | append(Rest, B)]
end.
-file("src/postgleam/pool.gleam", 113).
-spec handle_message(pool_state(), pool_message()) -> gleam@otp@actor:next(pool_state(), pool_message()).
handle_message(State, Msg) ->
case Msg of
{execute, Fun, Reply} ->
case erlang:element(2, State) of
[] ->
gleam@erlang@process:send(
Reply,
{error,
{connection_error,
<<"No connections available"/utf8>>}}
),
gleam@otp@actor:continue(State);
[Conn | Rest] ->
{Result, Conn@1} = Fun(
Conn,
erlang:element(4, State),
erlang:element(3, State)
),
gleam@erlang@process:send(Reply, Result),
gleam@otp@actor:continue(
{pool_state,
append(Rest, [Conn@1]),
erlang:element(3, State),
erlang:element(4, State),
erlang:element(5, State)}
)
end;
{simple_execute, Fun@1, Reply@1} ->
case erlang:element(2, State) of
[] ->
gleam@erlang@process:send(
Reply@1,
{error,
{connection_error,
<<"No connections available"/utf8>>}}
),
gleam@otp@actor:continue(State);
[Conn@2 | Rest@1] ->
{Result@1, Conn@3} = Fun@1(Conn@2, erlang:element(3, State)),
gleam@erlang@process:send(Reply@1, Result@1),
gleam@otp@actor:continue(
{pool_state,
append(Rest@1, [Conn@3]),
erlang:element(3, State),
erlang:element(4, State),
erlang:element(5, State)}
)
end;
{shutdown, Reply@2} ->
disconnect_all(erlang:element(2, State)),
gleam@erlang@process:send(Reply@2, nil),
gleam@otp@actor:stop()
end.
-file("src/postgleam/pool.gleam", 164).
-spec error_to_string(postgleam@error:error()) -> binary().
error_to_string(Err) ->
case Err of
{pg_error, Fields, _, _} ->
<<"PostgreSQL error: "/utf8, (erlang:element(4, Fields))/binary>>;
{connection_error, Msg} ->
Msg;
{authentication_error, Msg@1} ->
Msg@1;
{encode_error, Msg@2} ->
Msg@2;
{decode_error, Msg@3} ->
Msg@3;
{protocol_error, Msg@4} ->
Msg@4;
{socket_error, Msg@5} ->
Msg@5;
timeout_error ->
<<"Timeout"/utf8>>
end.
-file("src/postgleam/pool.gleam", 48).
?DOC(" Start a connection pool with the given config and size\n").
-spec start(postgleam@config:config(), integer()) -> {ok,
gleam@otp@actor:started(gleam@erlang@process:subject(pool_message()))} |
{error, binary()}.
start(Config, Size) ->
case begin
_pipe@2 = gleam@otp@actor:new_with_initialiser(
(erlang:element(8, Config) * Size) + 5000,
fun(Subject) -> case connect_pool(Config, Size, []) of
{ok, Conns} ->
Reg = postgleam@codec@registry:build(
postgleam@codec@defaults:matchers()
),
State = {pool_state, Conns, Config, Reg, Size},
_pipe = gleam@otp@actor:initialised(State),
_pipe@1 = gleam@otp@actor:returning(_pipe, Subject),
{ok, _pipe@1};
{error, E} ->
{error, error_to_string(E)}
end end
),
_pipe@3 = gleam@otp@actor:on_message(_pipe@2, fun handle_message/2),
gleam@otp@actor:start(_pipe@3)
end of
{ok, Started} ->
{ok, Started};
{error, init_timeout} ->
{error, <<"Pool initialization timed out"/utf8>>};
{error, {init_failed, Reason}} ->
{error, Reason};
{error, {init_exited, _}} ->
{error, <<"Pool process exited during init"/utf8>>}
end.
-file("src/postgleam/pool.gleam", 182).
?DOC(" Execute a parameterized query through the pool\n").
-spec 'query'(
gleam@erlang@process:subject(pool_message()),
binary(),
list(gleam@option:option(postgleam@value:value())),
integer()
) -> {ok, postgleam@connection:extended_query_result()} |
{error, postgleam@error:error()}.
'query'(Pool, Sql, Params, Timeout) ->
gleam@erlang@process:call(
Pool,
Timeout,
fun(Reply) ->
{execute,
fun(Conn, Reg, Config) ->
case postgleam@connection:extended_query(
Conn,
Sql,
Params,
Reg,
erlang:element(7, Config)
) of
{ok, {Result, Conn@1}} ->
{{ok, Result}, Conn@1};
{error, E} ->
{{error, E}, Conn}
end
end,
Reply}
end
).
-file("src/postgleam/pool.gleam", 202).
?DOC(" Execute a simple query through the pool\n").
-spec simple_query(
gleam@erlang@process:subject(pool_message()),
binary(),
integer()
) -> {ok, list(postgleam@connection:simple_query_result())} |
{error, postgleam@error:error()}.
simple_query(Pool, Sql, Timeout) ->
gleam@erlang@process:call(
Pool,
Timeout,
fun(Reply) ->
{simple_execute,
fun(Conn, Config) ->
case postgleam@connection:simple_query(
Conn,
Sql,
erlang:element(7, Config)
) of
{ok, {Results, Conn@1}} ->
{{ok, Results}, Conn@1};
{error, E} ->
{{error, E}, Conn}
end
end,
Reply}
end
).
-file("src/postgleam/pool.gleam", 249).
?DOC(" Shut down the pool, disconnecting all connections\n").
-spec shutdown(gleam@erlang@process:subject(pool_message()), integer()) -> nil.
shutdown(Pool, Timeout) ->
gleam@erlang@process:call(
Pool,
Timeout,
fun(Reply) -> {shutdown, Reply} end
).
-file("src/postgleam/pool.gleam", 253).
-spec decode_rows(
list(list(gleam@option:option(postgleam@value:value()))),
postgleam@decode:row_decoder(JGO),
list(JGO)
) -> {ok, list(JGO)} | {error, postgleam@error:error()}.
decode_rows(Rows, Decoder, Acc) ->
case Rows of
[] ->
{ok, lists:reverse(Acc)};
[Row | Rest] ->
case postgleam@decode:run(Decoder, Row) of
{ok, Val} ->
decode_rows(Rest, Decoder, [Val | Acc]);
{error, E} ->
{error, E}
end
end.
-file("src/postgleam/pool.gleam", 221).
?DOC(" Execute a parameterized query through the pool and decode rows.\n").
-spec query_with(
gleam@erlang@process:subject(pool_message()),
binary(),
list(gleam@option:option(postgleam@value:value())),
postgleam@decode:row_decoder(JGF),
integer()
) -> {ok, pool_response(JGF)} | {error, postgleam@error:error()}.
query_with(Pool, Sql, Params, Decoder, Timeout) ->
case 'query'(Pool, Sql, Params, Timeout) of
{ok, Result} ->
case decode_rows(erlang:element(4, Result), Decoder, []) of
{ok, Decoded} ->
{ok,
{pool_response,
Decoded,
erlang:length(Decoded),
erlang:element(2, Result)}};
{error, E} ->
{error, E}
end;
{error, E@1} ->
{error, E@1}
end.