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, connection_slot/0, queued_execute/0, queued_simple/0, queued_request/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)} |
health_check |
{reconnect, integer()}.
-type connection_slot() :: {active, postgleam@connection:connection_state()} |
{checked_out, postgleam@connection:connection_state()} |
{reconnecting, integer()}.
-type queued_execute() :: {queued_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()}),
integer()}.
-type queued_simple() :: {queued_simple,
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()}),
integer()}.
-type queued_request() :: {queued_exec, queued_execute()} |
{queued_simp, queued_simple()}.
-type pool_state() :: {pool_state,
list(connection_slot()),
postgleam@config:config(),
gleam@dict:dict(integer(), postgleam@codec:codec()),
integer(),
gleam@erlang@process:subject(pool_message()),
list(queued_request())}.
-type pool_response(KYC) :: {pool_response, list(KYC), integer(), binary()}.
-file("src/postgleam/pool.gleam", 159).
?DOC(
" Query pg_type for custom enum types and register their OIDs with\n"
" the text codec. Enums are text-representable on the wire — their\n"
" binary format is the enum label as UTF-8, identical to the text codec.\n"
).
-spec discover_enum_types(
list(postgleam@connection:connection_state()),
gleam@dict:dict(integer(), postgleam@codec:codec()),
postgleam@config:config()
) -> {gleam@dict:dict(integer(), postgleam@codec:codec()),
list(postgleam@connection:connection_state())}.
discover_enum_types(Conns, Reg, Config) ->
case Conns of
[First | Rest] ->
case postgleam@connection:simple_query(
First,
<<"SELECT oid::text FROM pg_type WHERE typtype = 'e'"/utf8>>,
erlang:element(7, Config)
) of
{ok, {Results, Updated_first}} ->
Reg@2 = case Results of
[Result | _] ->
gleam@list:fold(
erlang:element(4, Result),
Reg,
fun(Reg@1, Row) -> case Row of
[{some, Oid_str} | _] ->
case gleam_stdlib:parse_int(Oid_str) of
{ok, Oid} ->
postgleam@codec@registry:register(
Reg@1,
Oid,
postgleam@codec@text:matcher(
)
);
{error, _} ->
Reg@1
end;
_ ->
Reg@1
end end
);
_ ->
Reg
end,
{Reg@2, [Updated_first | Rest]};
{error, _} ->
{Reg, Conns}
end;
[] ->
{Reg, Conns}
end.
-file("src/postgleam/pool.gleam", 198).
-spec disconnect_all_conns(list(postgleam@connection:connection_state())) -> nil.
disconnect_all_conns(Conns) ->
case Conns of
[] ->
nil;
[Conn | Rest] ->
postgleam@connection:disconnect(Conn),
disconnect_all_conns(Rest)
end.
-file("src/postgleam/pool.gleam", 138).
-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_conns(Acc),
{error, E}
end
end.
-file("src/postgleam/pool.gleam", 334).
-spec find_available(list(connection_slot()), integer()) -> {ok,
{integer(), postgleam@connection:connection_state()}} |
{error, nil}.
find_available(Slots, Index) ->
case Slots of
[] ->
{error, nil};
[{active, Conn} | _] ->
{ok, {Index, Conn}};
[_ | Rest] ->
find_available(Rest, Index + 1)
end.
-file("src/postgleam/pool.gleam", 345).
-spec get_slot(list(connection_slot()), integer()) -> {ok, connection_slot()} |
{error, nil}.
get_slot(Slots, Index) ->
case {Slots, Index} of
{[], _} ->
{error, nil};
{[Slot | _], 0} ->
{ok, Slot};
{[_ | Rest], N} ->
get_slot(Rest, N - 1)
end.
-file("src/postgleam/pool.gleam", 353).
-spec set_slot(list(connection_slot()), integer(), connection_slot()) -> list(connection_slot()).
set_slot(Slots, Index, Value) ->
case {Slots, Index} of
{[], _} ->
[];
{[_ | Rest], 0} ->
[Value | Rest];
{[Slot | Rest@1], N} ->
[Slot | set_slot(Rest@1, N - 1, Value)]
end.
-file("src/postgleam/pool.gleam", 369).
-spec health_check_slots(
pool_state(),
list(connection_slot()),
integer(),
list(connection_slot())
) -> pool_state().
health_check_slots(State, Slots, Index, Acc) ->
case Slots of
[] ->
{pool_state,
lists:reverse(Acc),
erlang:element(3, State),
erlang:element(4, State),
erlang:element(5, State),
erlang:element(6, State),
erlang:element(7, State)};
[{active, Conn} | Rest] ->
case postgleam@connection:simple_query(
Conn,
<<"SELECT 1"/utf8>>,
erlang:element(7, erlang:element(3, State))
) of
{ok, {_, Conn@1}} ->
health_check_slots(
State,
Rest,
Index + 1,
[{active, Conn@1} | Acc]
);
{error, _} ->
_ = gleam@erlang@process:send_after(
erlang:element(6, State),
0,
{reconnect, Index}
),
health_check_slots(
State,
Rest,
Index + 1,
[{reconnecting, 0} | Acc]
)
end;
[Slot | Rest@1] ->
health_check_slots(State, Rest@1, Index + 1, [Slot | Acc])
end.
-file("src/postgleam/pool.gleam", 404).
-spec expire_queue(pool_state()) -> pool_state().
expire_queue(State) ->
Now = postgleam_ffi:monotonic_time_ms(),
Timeout = erlang:element(13, erlang:element(3, State)),
{Alive, Expired} = gleam@list:partition(
erlang:element(7, State),
fun(Req) ->
Enqueued = case Req of
{queued_exec, R} ->
erlang:element(4, R);
{queued_simp, R@1} ->
erlang:element(4, R@1)
end,
(Now - Enqueued) < Timeout
end
),
gleam@list:each(Expired, fun(Req@1) -> case Req@1 of
{queued_exec, R@2} ->
gleam@erlang@process:send(
erlang:element(3, R@2),
{error, timeout_error}
);
{queued_simp, R@3} ->
gleam@erlang@process:send(
erlang:element(3, R@3),
{error, timeout_error}
)
end end),
{pool_state,
erlang:element(2, State),
erlang:element(3, State),
erlang:element(4, State),
erlang:element(5, State),
erlang:element(6, State),
Alive}.
-file("src/postgleam/pool.gleam", 499).
-spec reject_queue(list(queued_request())) -> nil.
reject_queue(Queue) ->
case Queue of
[] ->
nil;
[{queued_exec, R} | Rest] ->
gleam@erlang@process:send(
erlang:element(3, R),
{error, {connection_error, <<"Pool shutting down"/utf8>>}}
),
reject_queue(Rest);
[{queued_simp, R@1} | Rest@1] ->
gleam@erlang@process:send(
erlang:element(3, R@1),
{error, {connection_error, <<"Pool shutting down"/utf8>>}}
),
reject_queue(Rest@1)
end.
-file("src/postgleam/pool.gleam", 513).
-spec disconnect_all_slots(list(connection_slot())) -> nil.
disconnect_all_slots(Slots) ->
case Slots of
[] ->
nil;
[{active, Conn} | Rest] ->
postgleam@connection:disconnect(Conn),
disconnect_all_slots(Rest);
[{checked_out, Conn@1} | Rest@1] ->
postgleam@connection:disconnect(Conn@1),
disconnect_all_slots(Rest@1);
[{reconnecting, _} | Rest@2] ->
disconnect_all_slots(Rest@2)
end.
-file("src/postgleam/pool.gleam", 532).
-spec is_socket_error(
{ok, postgleam@connection:extended_query_result()} |
{error, postgleam@error:error()}
) -> boolean().
is_socket_error(Result) ->
case Result of
{error, {socket_error, _}} ->
true;
_ ->
false
end.
-file("src/postgleam/pool.gleam", 539).
-spec is_simple_socket_error(
{ok, list(postgleam@connection:simple_query_result())} |
{error, postgleam@error:error()}
) -> boolean().
is_simple_socket_error(Result) ->
case Result of
{error, {socket_error, _}} ->
true;
_ ->
false
end.
-file("src/postgleam/pool.gleam", 425).
-spec drain_queue(pool_state()) -> pool_state().
drain_queue(State) ->
case erlang:element(7, State) of
[] ->
State;
[First | Rest] ->
Now = postgleam_ffi:monotonic_time_ms(),
case First of
{queued_exec, Req} ->
case (Now - erlang:element(4, Req)) >= erlang:element(
13,
erlang:element(3, State)
) of
true ->
gleam@erlang@process:send(
erlang:element(3, Req),
{error, timeout_error}
),
drain_queue(
{pool_state,
erlang:element(2, State),
erlang:element(3, State),
erlang:element(4, State),
erlang:element(5, State),
erlang:element(6, State),
Rest}
);
false ->
case find_available(erlang:element(2, State), 0) of
{ok, {Index, Conn}} ->
{Result, Conn@1} = (erlang:element(2, Req))(
Conn,
erlang:element(4, State),
erlang:element(3, State)
),
Is_dead = is_socket_error(Result),
gleam@erlang@process:send(
erlang:element(3, Req),
Result
),
case Is_dead of
true ->
Slots = set_slot(
erlang:element(2, State),
Index,
{reconnecting, 0}
),
_ = gleam@erlang@process:send_after(
erlang:element(6, State),
0,
{reconnect, Index}
),
drain_queue(
{pool_state,
Slots,
erlang:element(3, State),
erlang:element(4, State),
erlang:element(5, State),
erlang:element(6, State),
Rest}
);
false ->
Slots@1 = set_slot(
erlang:element(2, State),
Index,
{active, Conn@1}
),
drain_queue(
{pool_state,
Slots@1,
erlang:element(3, State),
erlang:element(4, State),
erlang:element(5, State),
erlang:element(6, State),
Rest}
)
end;
{error, _} ->
State
end
end;
{queued_simp, Req@1} ->
case (Now - erlang:element(4, Req@1)) >= erlang:element(
13,
erlang:element(3, State)
) of
true ->
gleam@erlang@process:send(
erlang:element(3, Req@1),
{error, timeout_error}
),
drain_queue(
{pool_state,
erlang:element(2, State),
erlang:element(3, State),
erlang:element(4, State),
erlang:element(5, State),
erlang:element(6, State),
Rest}
);
false ->
case find_available(erlang:element(2, State), 0) of
{ok, {Index@1, Conn@2}} ->
{Result@1, Conn@3} = (erlang:element(
2,
Req@1
))(Conn@2, erlang:element(3, State)),
Is_dead@1 = is_simple_socket_error(Result@1),
gleam@erlang@process:send(
erlang:element(3, Req@1),
Result@1
),
case Is_dead@1 of
true ->
Slots@2 = set_slot(
erlang:element(2, State),
Index@1,
{reconnecting, 0}
),
_ = gleam@erlang@process:send_after(
erlang:element(6, State),
0,
{reconnect, Index@1}
),
drain_queue(
{pool_state,
Slots@2,
erlang:element(3, State),
erlang:element(4, State),
erlang:element(5, State),
erlang:element(6, State),
Rest}
);
false ->
Slots@3 = set_slot(
erlang:element(2, State),
Index@1,
{active, Conn@3}
),
drain_queue(
{pool_state,
Slots@3,
erlang:element(3, State),
erlang:element(4, State),
erlang:element(5, State),
erlang:element(6, State),
Rest}
)
end;
{error, _} ->
State
end
end
end
end.
-file("src/postgleam/pool.gleam", 552).
-spec pow2(integer()) -> integer().
pow2(N) ->
case N of
0 ->
1;
_ ->
2 * pow2(N - 1)
end.
-file("src/postgleam/pool.gleam", 559).
-spec append(list(KZH), list(KZH)) -> list(KZH).
append(A, B) ->
case A of
[] ->
B;
[X | Rest] ->
[X | append(Rest, B)]
end.
-file("src/postgleam/pool.gleam", 208).
-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 find_available(erlang:element(2, State), 0) of
{ok, {Index, Conn}} ->
{Result, Conn@1} = Fun(
Conn,
erlang:element(4, State),
erlang:element(3, State)
),
Is_dead = is_socket_error(Result),
gleam@erlang@process:send(Reply, Result),
case Is_dead of
true ->
Slots = set_slot(
erlang:element(2, State),
Index,
{reconnecting, 0}
),
_ = gleam@erlang@process:send_after(
erlang:element(6, State),
0,
{reconnect, Index}
),
gleam@otp@actor:continue(
{pool_state,
Slots,
erlang:element(3, State),
erlang:element(4, State),
erlang:element(5, State),
erlang:element(6, State),
erlang:element(7, State)}
);
false ->
Slots@1 = set_slot(
erlang:element(2, State),
Index,
{active, Conn@1}
),
gleam@otp@actor:continue(
{pool_state,
Slots@1,
erlang:element(3, State),
erlang:element(4, State),
erlang:element(5, State),
erlang:element(6, State),
erlang:element(7, State)}
)
end;
{error, _} ->
Now = postgleam_ffi:monotonic_time_ms(),
Req = {queued_exec, {queued_execute, Fun, Reply, Now}},
gleam@otp@actor:continue(
{pool_state,
erlang:element(2, State),
erlang:element(3, State),
erlang:element(4, State),
erlang:element(5, State),
erlang:element(6, State),
append(erlang:element(7, State), [Req])}
)
end;
{simple_execute, Fun@1, Reply@1} ->
case find_available(erlang:element(2, State), 0) of
{ok, {Index@1, Conn@2}} ->
{Result@1, Conn@3} = Fun@1(Conn@2, erlang:element(3, State)),
Is_dead@1 = is_simple_socket_error(Result@1),
gleam@erlang@process:send(Reply@1, Result@1),
case Is_dead@1 of
true ->
Slots@2 = set_slot(
erlang:element(2, State),
Index@1,
{reconnecting, 0}
),
_ = gleam@erlang@process:send_after(
erlang:element(6, State),
0,
{reconnect, Index@1}
),
gleam@otp@actor:continue(
{pool_state,
Slots@2,
erlang:element(3, State),
erlang:element(4, State),
erlang:element(5, State),
erlang:element(6, State),
erlang:element(7, State)}
);
false ->
Slots@3 = set_slot(
erlang:element(2, State),
Index@1,
{active, Conn@3}
),
gleam@otp@actor:continue(
{pool_state,
Slots@3,
erlang:element(3, State),
erlang:element(4, State),
erlang:element(5, State),
erlang:element(6, State),
erlang:element(7, State)}
)
end;
{error, _} ->
Now@1 = postgleam_ffi:monotonic_time_ms(),
Req@1 = {queued_simp,
{queued_simple, Fun@1, Reply@1, Now@1}},
gleam@otp@actor:continue(
{pool_state,
erlang:element(2, State),
erlang:element(3, State),
erlang:element(4, State),
erlang:element(5, State),
erlang:element(6, State),
append(erlang:element(7, State), [Req@1])}
)
end;
health_check ->
State@1 = health_check_slots(State, erlang:element(2, State), 0, []),
State@2 = expire_queue(State@1),
_ = gleam@erlang@process:send_after(
erlang:element(6, State@2),
erlang:element(12, erlang:element(3, State@2)),
health_check
),
gleam@otp@actor:continue(State@2);
{reconnect, Index@2} ->
case get_slot(erlang:element(2, State), Index@2) of
{ok, {reconnecting, Attempts}} ->
case postgleam@connection:connect(erlang:element(3, State)) of
{ok, Conn@4} ->
State@3 = {pool_state,
set_slot(
erlang:element(2, State),
Index@2,
{active, Conn@4}
),
erlang:element(3, State),
erlang:element(4, State),
erlang:element(5, State),
erlang:element(6, State),
erlang:element(7, State)},
State@4 = drain_queue(State@3),
gleam@otp@actor:continue(State@4);
{error, _} ->
Delay = gleam@int:min(30000, pow2(Attempts) * 1000),
_ = gleam@erlang@process:send_after(
erlang:element(6, State),
Delay,
{reconnect, Index@2}
),
Slots@4 = set_slot(
erlang:element(2, State),
Index@2,
{reconnecting, Attempts + 1}
),
gleam@otp@actor:continue(
{pool_state,
Slots@4,
erlang:element(3, State),
erlang:element(4, State),
erlang:element(5, State),
erlang:element(6, State),
erlang:element(7, State)}
)
end;
_ ->
gleam@otp@actor:continue(State)
end;
{shutdown, Reply@2} ->
reject_queue(erlang:element(7, State)),
disconnect_all_slots(erlang:element(2, State)),
gleam@erlang@process:send(Reply@2, nil),
gleam@otp@actor:stop()
end.
-file("src/postgleam/pool.gleam", 566).
-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", 94).
?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@defaults:build_registry(),
{Reg@1, Conns@1} = discover_enum_types(
Conns,
Reg,
Config
),
Slots = gleam@list:map(
Conns@1,
fun(C) -> {active, C} end
),
State = {pool_state,
Slots,
Config,
Reg@1,
Size,
Subject,
[]},
_ = gleam@erlang@process:send_after(
Subject,
erlang:element(12, Config),
health_check
),
_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", 585).
?DOC(
" Execute a parameterized query through the pool.\n"
" Uses PgBouncer-safe unnamed statements when pgbouncer mode is enabled.\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) ->
Result = case erlang:element(11, Config) of
true ->
postgleam@connection:extended_query_unnamed(
Conn,
Sql,
Params,
Reg,
erlang:element(7, Config)
);
false ->
postgleam@connection:extended_query(
Conn,
Sql,
Params,
Reg,
erlang:element(7, Config)
)
end,
case Result of
{ok, {R, C}} ->
{{ok, R}, C};
{error, E} ->
{{error, E}, Conn}
end
end,
Reply}
end
).
-file("src/postgleam/pool.gleam", 613).
?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", 660).
?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", 664).
-spec decode_rows(
list(list(gleam@option:option(postgleam@value:value()))),
postgleam@decode:row_decoder(LAG),
list(LAG)
) -> {ok, list(LAG)} | {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", 632).
?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(KZX),
integer()
) -> {ok, pool_response(KZX)} | {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.