Current section
Files
Jump to
Current section
Files
src/pgl@internal@protocol.erl
-module(pgl@internal@protocol).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch, inline]).
-define(FILEPATH, "src/pgl/internal/protocol.gleam").
-export([application/2, connection_parameters/2, username/2, password/2, database/2, ssl/2, auth/2, simple/2, extended/0, on_param_description/2, on_decode_row/2, pipeline/0, process/3, batch_process/4]).
-export_type([config/0, extended/1, pipeline/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.
?MODULEDOC(false).
-opaque config() :: {config,
binary(),
binary(),
gleam@option:option(binary()),
binary(),
list({binary(), binary()}),
pgl@internal:ssl()}.
-type extended(MHU) :: {extended,
boolean(),
fun((list(gleam@option:option(bitstring())), list(integer())) -> {ok,
list(gleam@dynamic:dynamic_())} |
{error, pgl@internal:internal_error()}),
fun((binary(), list(MHU), list(integer())) -> {ok, bitstring()} |
{error, pgl@internal:internal_error()}),
list(pgl@internal:row_description_field()),
list(binary()),
list(list(gleam@dynamic:dynamic_())),
integer()}.
-type pipeline(MHV) :: {pipeline, integer(), integer(), list(extended(MHV))}.
-file("src/pgl/internal/protocol.gleam", 38).
?DOC(false).
-spec application(config(), binary()) -> config().
application(Conf, Application) ->
{config,
erlang:element(2, Conf),
erlang:element(3, Conf),
erlang:element(4, Conf),
Application,
erlang:element(6, Conf),
erlang:element(7, Conf)}.
-file("src/pgl/internal/protocol.gleam", 42).
?DOC(false).
-spec connection_parameters(config(), list({binary(), binary()})) -> config().
connection_parameters(Conf, Connection_parameters) ->
{config,
erlang:element(2, Conf),
erlang:element(3, Conf),
erlang:element(4, Conf),
erlang:element(5, Conf),
Connection_parameters,
erlang:element(7, Conf)}.
-file("src/pgl/internal/protocol.gleam", 49).
?DOC(false).
-spec username(config(), binary()) -> config().
username(Conf, Username) ->
{config,
erlang:element(2, Conf),
Username,
erlang:element(4, Conf),
erlang:element(5, Conf),
erlang:element(6, Conf),
erlang:element(7, Conf)}.
-file("src/pgl/internal/protocol.gleam", 53).
?DOC(false).
-spec password(config(), binary()) -> config().
password(Conf, Password) ->
{config,
erlang:element(2, Conf),
erlang:element(3, Conf),
{some, Password},
erlang:element(5, Conf),
erlang:element(6, Conf),
erlang:element(7, Conf)}.
-file("src/pgl/internal/protocol.gleam", 57).
?DOC(false).
-spec database(config(), binary()) -> config().
database(Conf, Database) ->
{config,
Database,
erlang:element(3, Conf),
erlang:element(4, Conf),
erlang:element(5, Conf),
erlang:element(6, Conf),
erlang:element(7, Conf)}.
-file("src/pgl/internal/protocol.gleam", 61).
?DOC(false).
-spec ssl(config(), pgl@internal:ssl()) -> config().
ssl(Conf, Ssl) ->
{config,
erlang:element(2, Conf),
erlang:element(3, Conf),
erlang:element(4, Conf),
erlang:element(5, Conf),
erlang:element(6, Conf),
Ssl}.
-file("src/pgl/internal/protocol.gleam", 239).
?DOC(false).
-spec handle_error_response(gleam@dict:dict(bitstring(), binary())) -> {ok,
any()} |
{error, pgl@internal:internal_error()}.
handle_error_response(Fields) ->
Code = begin
_pipe = gleam_stdlib:map_get(Fields, <<"C"/utf8>>),
gleam@result:unwrap(_pipe, <<""/utf8>>)
end,
Message = begin
_pipe@1 = gleam_stdlib:map_get(Fields, <<"M"/utf8>>),
gleam@result:unwrap(_pipe@1, <<""/utf8>>)
end,
Name = begin
_pipe@2 = pgl@internal:pg_error_code_name(Code),
gleam@result:unwrap(_pipe@2, <<""/utf8>>)
end,
_pipe@3 = {postgres_error, Code, Name, Message, Fields},
{error, _pipe@3}.
-file("src/pgl/internal/protocol.gleam", 286).
?DOC(false).
-spec auth_sasl_final(bitstring(), bitstring()) -> {ok, bitstring()} |
{error, pgl@internal:internal_error()}.
auth_sasl_final(Server_final, Server_signature) ->
gleam@result:'try'(
pgl@internal@scram:parse_server_final(Server_final),
fun(Srv_final) ->
case gleam@crypto:secure_compare(Srv_final, Server_signature) of
true ->
{ok, Server_signature};
false ->
_pipe = {authentication_error,
authentication_failed,
<<"Failed to match server signature"/utf8>>},
{error, _pipe}
end
end
).
-file("src/pgl/internal/protocol.gleam", 250).
?DOC(false).
-spec auth_sasl_continue(
pgl@internal@socket:socket(),
config(),
bitstring(),
bitstring()
) -> {ok, bitstring()} | {error, pgl@internal:internal_error()}.
auth_sasl_continue(Sock, Conf, Server_first, Client_nonce) ->
_pipe = pgl@internal@scram:parse_server_first(Server_first, Client_nonce),
gleam@result:'try'(
_pipe,
fun(Sf) ->
User = <<(erlang:element(3, Conf))/binary>>,
case erlang:element(4, Conf) of
none ->
_pipe@1 = authentication_failed,
_pipe@2 = {authentication_error,
_pipe@1,
<<"Server requested SCRAM authentication but no password was provided"/utf8>>},
{error, _pipe@2};
{some, Password} ->
Pass = <<Password/binary>>,
gleam@result:'try'(
pgl@internal@scram:client_final(
Sf,
Client_nonce,
User,
Pass
),
fun(_use0) ->
{Client_final, Server_signature} = _use0,
Encoded_client_final = pgl@internal@encode:scram_response(
Client_final
),
_pipe@3 = pgl@internal@socket:send(
Sock,
Encoded_client_final
),
gleam@result:replace(_pipe@3, Server_signature)
end
)
end
end
).
-file("src/pgl/internal/protocol.gleam", 212).
?DOC(false).
-spec auth_sasl(pgl@internal@socket:socket(), list(binary()), config()) -> {ok,
bitstring()} |
{error, pgl@internal:internal_error()}.
auth_sasl(Sock, Methods, Conf) ->
case Methods of
[<<"SCRAM-SHA-256"/utf8>>] ->
Client_nonce = pgl@internal@scram:get_nonce(16),
_pipe = pgl@internal@scram:client_first(
<<(erlang:element(3, Conf))/binary>>,
Client_nonce
),
gleam@result:'try'(
_pipe,
fun(Client_first) -> _pipe@1 = Client_first,
_pipe@2 = pgl@internal@encode:auth_scram_client_first(
_pipe@1
),
_pipe@3 = pgl@internal@socket:send(Sock, _pipe@2),
gleam@result:replace(_pipe@3, Client_nonce) end
);
_ ->
_pipe@4 = {authentication_error,
method_not_implemented,
<<"Supported methods: [SCRAM-SHA-256]"/utf8>>},
{error, _pipe@4}
end.
-file("src/pgl/internal/protocol.gleam", 173).
?DOC(false).
-spec md5_password(binary(), binary(), bitstring()) -> binary().
md5_password(Username, Password, Salt) ->
Inner = begin
_pipe = gleam@crypto:hash(md5, <<Password/binary, Username/binary>>),
_pipe@1 = gleam_stdlib:base16_encode(_pipe),
string:lowercase(_pipe@1)
end,
Outer = begin
_pipe@2 = gleam@crypto:hash(md5, <<Inner/binary, Salt/bitstring>>),
_pipe@3 = gleam_stdlib:base16_encode(_pipe@2),
string:lowercase(_pipe@3)
end,
<<"md5"/utf8, Outer/binary>>.
-file("src/pgl/internal/protocol.gleam", 469).
?DOC(false).
-spec receive_message(pgl@internal@socket:socket()) -> {ok,
pgl@internal:message()} |
{error, pgl@internal:internal_error()}.
receive_message(Sock) ->
gleam@result:'try'(
pgl@internal@socket:'receive'(Sock, 5),
fun(Data) -> case Data of
<<Code:8/bitstring, Size:32/integer>> ->
case Size - 4 of
0 ->
pgl@internal@decode:message(Code, <<>>);
Size1 ->
gleam@result:'try'(
pgl@internal@socket:'receive'(Sock, Size1),
fun(Payload) ->
pgl@internal@decode:message(Code, Payload)
end
)
end;
_ ->
_pipe = decoding_error,
_pipe@1 = {protocol_error,
_pipe,
<<"Unexpected data received"/utf8>>},
{error, _pipe@1}
end end
).
-file("src/pgl/internal/protocol.gleam", 187).
?DOC(false).
-spec do_password_auth(
config(),
fun((binary()) -> binary()),
binary(),
pgl@internal@socket:socket()
) -> {ok, pgl@internal@socket:socket()} | {error, pgl@internal:internal_error()}.
do_password_auth(Conf, Process_password, Kind, Sock) ->
case erlang:element(4, Conf) of
{some, Pass} ->
_pipe = Pass,
_pipe@1 = Process_password(_pipe),
_pipe@2 = pgl@internal@encode:password(_pipe@1),
_pipe@3 = pgl@internal@socket:send(Sock, _pipe@2),
gleam@result:'try'(
_pipe@3,
fun(_capture) -> auth_flow(_capture, Conf, <<>>) end
);
none ->
_pipe@4 = authentication_failed,
_pipe@5 = {authentication_error,
_pipe@4,
<<<<"Server requested "/utf8, Kind/binary>>/binary,
"authentication but no password was provided"/utf8>>},
{error, _pipe@5}
end.
-file("src/pgl/internal/protocol.gleam", 127).
?DOC(false).
-spec auth_flow(pgl@internal@socket:socket(), config(), bitstring()) -> {ok,
pgl@internal@socket:socket()} |
{error, pgl@internal:internal_error()}.
auth_flow(Sock, Conf, Prev) ->
gleam@result:'try'(receive_message(Sock), fun(Msg) -> case Msg of
authentication_ok ->
auth_flow(Sock, Conf, Prev);
{authentication_m_d5_password, Salt} ->
do_password_auth(
Conf,
fun(_capture) ->
md5_password(
erlang:element(3, Conf),
_capture,
Salt
)
end,
<<"MD5"/utf8>>,
Sock
);
authentication_cleartext_password ->
do_password_auth(
Conf,
fun gleam@function:identity/1,
<<"password"/utf8>>,
Sock
);
{authentication_s_a_s_l, Methods} ->
gleam@result:'try'(
auth_sasl(Sock, Methods, Conf),
fun(Nonce) -> auth_flow(Sock, Conf, Nonce) end
);
{authentication_s_a_s_l_continue, First} ->
gleam@result:'try'(
auth_sasl_continue(Sock, Conf, First, Prev),
fun(Srv_sig) -> auth_flow(Sock, Conf, Srv_sig) end
);
{authentication_s_a_s_l_final, Server_final} ->
gleam@result:'try'(
auth_sasl_final(Server_final, Prev),
fun(_) -> auth_flow(Sock, Conf, <<>>) end
);
{error_response, Fields} ->
handle_error_response(Fields);
{backend_key_data, _, _} ->
auth_flow(Sock, Conf, <<>>);
{notification_response, _, _, _} ->
auth_flow(Sock, Conf, <<>>);
{notice_response, _} ->
auth_flow(Sock, Conf, <<>>);
{parameter_status, Name, Value} ->
_pipe = Sock,
_pipe@1 = pgl@internal@socket:parameter(_pipe, Name, Value),
auth_flow(_pipe@1, Conf, <<>>);
{ready_for_query, _} ->
{ok, Sock};
_ ->
_pipe@2 = message_error,
_pipe@3 = {protocol_error,
_pipe@2,
<<"Unexpected message"/utf8>>},
{error, _pipe@3}
end end).
-file("src/pgl/internal/protocol.gleam", 111).
?DOC(false).
-spec setup(pgl@internal@socket:socket(), config()) -> {ok,
pgl@internal@socket:socket()} |
{error, pgl@internal:internal_error()}.
setup(Sock, Conf) ->
Message = begin
_pipe = [{<<"user"/utf8>>, erlang:element(3, Conf)},
{<<"database"/utf8>>, erlang:element(2, Conf)},
{<<"application_name"/utf8>>, erlang:element(5, Conf)} |
erlang:element(6, Conf)],
pgl@internal@encode:startup(_pipe)
end,
gleam@result:'try'(
pgl@internal@socket:send(Sock, Message),
fun(Sock@1) -> auth_flow(Sock@1, Conf, <<>>) end
).
-file("src/pgl/internal/protocol.gleam", 89).
?DOC(false).
-spec do_ssl_upgrade(pgl@internal@socket:socket(), boolean()) -> {ok,
pgl@internal@socket:socket()} |
{error, pgl@internal:internal_error()}.
do_ssl_upgrade(Sock, Verified) ->
gleam@result:'try'(
pgl@internal@socket:send(Sock, pgl@internal@encode:ssl_request()),
fun(Sock@1) -> case pgl@internal@socket:'receive'(Sock@1, 1) of
{ok, <<"S"/utf8>>} ->
pgl@internal@socket:to_ssl(Sock@1, Verified);
{ok, <<"N"/utf8>>} ->
_pipe = ssl_error,
_pipe@1 = {protocol_error, _pipe, <<"SSL Refused"/utf8>>},
{error, _pipe@1};
{ok, _} ->
_pipe@2 = ssl_error,
_pipe@3 = {protocol_error,
_pipe@2,
<<"Failed to upgrade SSL"/utf8>>},
{error, _pipe@3};
{error, Err} ->
{error, Err}
end end
).
-file("src/pgl/internal/protocol.gleam", 78).
?DOC(false).
-spec ssl_upgrade(pgl@internal@socket:socket(), pgl@internal:ssl()) -> {ok,
pgl@internal@socket:socket()} |
{error, pgl@internal:internal_error()}.
ssl_upgrade(Sock, Ssl) ->
case Ssl of
ssl_verified ->
do_ssl_upgrade(Sock, true);
ssl_unverified ->
do_ssl_upgrade(Sock, false);
ssl_disabled ->
{ok, Sock}
end.
-file("src/pgl/internal/protocol.gleam", 67).
?DOC(false).
-spec auth(pgl@internal@socket:socket(), config()) -> {ok,
pgl@internal@socket:socket()} |
{error, pgl@internal:internal_error()}.
auth(Sock, Conf) ->
_pipe = Sock,
_pipe@1 = ssl_upgrade(_pipe, erlang:element(7, Conf)),
gleam@result:'try'(_pipe@1, fun(_capture) -> setup(_capture, Conf) end).
-file("src/pgl/internal/protocol.gleam", 320).
?DOC(false).
-spec simple_flow(
pgl@internal@socket:socket(),
list(list(gleam@option:option(bitstring())))
) -> {ok, list(list(gleam@option:option(bitstring())))} |
{error, pgl@internal:internal_error()}.
simple_flow(Sock, Acc) ->
gleam@result:'try'(receive_message(Sock), fun(Msg) -> case Msg of
{command_complete, _, _} ->
simple_flow(Sock, Acc);
{data_row, Values} ->
simple_flow(Sock, [Values | Acc]);
{error_response, Fields} ->
handle_error_response(Fields);
{notice_response, _} ->
simple_flow(Sock, Acc);
{notification_response, _, _, _} ->
simple_flow(Sock, Acc);
{ready_for_query, _} ->
{ok, Acc};
{row_description, _, _} ->
simple_flow(Sock, Acc);
_ ->
_pipe = {protocol_error,
message_error,
<<"Unexpected message in simple flow"/utf8>>},
{error, _pipe}
end end).
-file("src/pgl/internal/protocol.gleam", 311).
?DOC(false).
-spec simple(bitstring(), pgl@internal@socket:socket()) -> {ok,
list(list(gleam@option:option(bitstring())))} |
{error, pgl@internal:internal_error()}.
simple(Packet, Sock) ->
gleam@result:'try'(
pgl@internal@socket:send(Sock, Packet),
fun(Sock@1) -> simple_flow(Sock@1, []) end
).
-file("src/pgl/internal/protocol.gleam", 344).
?DOC(false).
-spec flush(
{ok, MJN} | {error, pgl@internal:internal_error()},
pgl@internal@socket:socket()
) -> {ok, MJN} | {error, pgl@internal:internal_error()}.
flush(Res, Sock) ->
gleam@result:'try'(receive_message(Sock), fun(Msg) -> case Msg of
{parameter_status, _, _} ->
flush(Res, Sock);
{ready_for_query, _} ->
Res;
_ ->
flush(Res, Sock)
end end).
-file("src/pgl/internal/protocol.gleam", 357).
?DOC(false).
-spec sync(pgl@internal@socket:socket()) -> {ok, pgl@internal@socket:socket()} |
{error, pgl@internal:internal_error()}.
sync(Sock) ->
_pipe = pgl@internal@encode:sync(),
_pipe@1 = pgl@internal@socket:send(Sock, _pipe),
_pipe@2 = gleam@result:'try'(_pipe@1, fun receive_message/1),
gleam@result:'try'(_pipe@2, fun(Msg) -> case Msg of
{ready_for_query, _} ->
{ok, Sock};
_ ->
_pipe@3 = message_error,
_pipe@4 = {protocol_error,
_pipe@3,
<<"Expected ReadyForQuery after Sync"/utf8>>},
{error, _pipe@4}
end end).
-file("src/pgl/internal/protocol.gleam", 394).
?DOC(false).
-spec extended() -> extended(any()).
extended() ->
{extended,
false,
fun(_, _) -> {ok, []} end,
fun(_, _, _) -> {ok, <<>>} end,
[],
[],
[],
0}.
-file("src/pgl/internal/protocol.gleam", 406).
?DOC(false).
-spec on_param_description(
extended(MJW),
fun((binary(), list(MJW), list(integer())) -> {ok, bitstring()} |
{error, pgl@internal:internal_error()})
) -> extended(MJW).
on_param_description(Ext, Handle_param_description) ->
{extended,
erlang:element(2, Ext),
erlang:element(3, Ext),
Handle_param_description,
erlang:element(5, Ext),
erlang:element(6, Ext),
erlang:element(7, Ext),
erlang:element(8, Ext)}.
-file("src/pgl/internal/protocol.gleam", 413).
?DOC(false).
-spec on_decode_row(
extended(MKA),
fun((list(gleam@option:option(bitstring())), list(integer())) -> {ok,
list(gleam@dynamic:dynamic_())} |
{error, pgl@internal:internal_error()})
) -> extended(MKA).
on_decode_row(Ext, Handle_decode_row) ->
{extended,
erlang:element(2, Ext),
Handle_decode_row,
erlang:element(4, Ext),
erlang:element(5, Ext),
erlang:element(6, Ext),
erlang:element(7, Ext),
erlang:element(8, Ext)}.
-file("src/pgl/internal/protocol.gleam", 446).
?DOC(false).
-spec handle_row_description(
extended(MKL),
list(pgl@internal:row_description_field())
) -> extended(MKL).
handle_row_description(Ext, Descriptions) ->
Fields = gleam@list:map(
Descriptions,
fun(Desc) -> erlang:element(2, Desc) end
),
{extended,
erlang:element(2, Ext),
erlang:element(3, Ext),
erlang:element(4, Ext),
Descriptions,
Fields,
erlang:element(7, Ext),
erlang:element(8, Ext)}.
-file("src/pgl/internal/protocol.gleam", 503).
?DOC(false).
-spec reverse_acc(pipeline(MLB)) -> pipeline(MLB).
reverse_acc(Pl) ->
{pipeline,
erlang:element(2, Pl),
erlang:element(3, Pl),
lists:reverse(erlang:element(4, Pl))}.
-file("src/pgl/internal/protocol.gleam", 511).
?DOC(false).
-spec increment_ready(pipeline(MLH)) -> pipeline(MLH).
increment_ready(Pl) ->
{pipeline,
erlang:element(2, Pl),
erlang:element(3, Pl) + 1,
erlang:element(4, Pl)}.
-file("src/pgl/internal/protocol.gleam", 507).
?DOC(false).
-spec increment_sync(pipeline(MLE)) -> pipeline(MLE).
increment_sync(Pl) ->
{pipeline,
erlang:element(2, Pl) + 1,
erlang:element(3, Pl),
erlang:element(4, Pl)}.
-file("src/pgl/internal/protocol.gleam", 610).
?DOC(false).
-spec error_response_cleanup(
{ok, MMH} | {error, pgl@internal:internal_error()},
boolean(),
integer(),
integer(),
pgl@internal@socket:socket()
) -> {ok, MMH} | {error, pgl@internal:internal_error()}.
error_response_cleanup(Err, Needs_sync, Syncs, Ready, Sock) ->
{Err@2, Syncs@1} = case Needs_sync of
false ->
{flush(Err, Sock), Syncs};
true ->
Err@1 = begin
_pipe = pgl@internal@encode:sync(),
_pipe@1 = pgl@internal@socket:send(Sock, _pipe),
gleam@result:'try'(_pipe@1, fun(_) -> flush(Err, Sock) end)
end,
{Err@1, Syncs + 1}
end,
Ready@1 = Ready + 1,
case Syncs@1 > Ready@1 of
true ->
error_response_cleanup(Err@2, false, Syncs@1, Ready@1, Sock);
false ->
Err@2
end.
-file("src/pgl/internal/protocol.gleam", 455).
?DOC(false).
-spec handle_data_row(
list(gleam@option:option(bitstring())),
extended(MKP),
fun((list(gleam@option:option(bitstring())), list(integer())) -> {ok,
list(gleam@dynamic:dynamic_())} |
{error, pgl@internal:internal_error()})
) -> {ok, extended(MKP)} | {error, pgl@internal:internal_error()}.
handle_data_row(Row, Rows, Decode_row) ->
Oids = gleam@list:map(
erlang:element(5, Rows),
fun(D) -> erlang:element(5, D) end
),
gleam@result:map(
Decode_row(Row, Oids),
fun(Values) ->
Values@1 = gleam@list:prepend(erlang:element(7, Rows), Values),
{extended,
erlang:element(2, Rows),
erlang:element(3, Rows),
erlang:element(4, Rows),
erlang:element(5, Rows),
erlang:element(6, Rows),
Values@1,
erlang:element(8, Rows)}
end
).
-file("src/pgl/internal/protocol.gleam", 499).
?DOC(false).
-spec set_acc(pipeline(MKW), list(extended(MKW))) -> pipeline(MKW).
set_acc(Pl, Acc) ->
{pipeline, erlang:element(2, Pl), erlang:element(3, Pl), Acc}.
-file("src/pgl/internal/protocol.gleam", 652).
?DOC(false).
-spec next_param_description(
pipeline(MMV),
pgl@internal@encode:'query'(MMV, MMX),
list(pgl@internal@encode:'query'(MMV, MMX)),
extended(MMV),
list(integer()),
pgl@internal@socket:socket()
) -> {ok, pipeline(MMV)} | {error, pgl@internal:internal_error()}.
next_param_description(Pl, Query, Rest, Ext, Oids, Sock) ->
Sql = erlang:element(2, Query),
Params = erlang:element(3, Query),
gleam@result:'try'(
(erlang:element(4, Ext))(Sql, Params, Oids),
fun(Packet) ->
gleam@result:'try'(
pgl@internal@socket:send(Sock, Packet),
fun(Sock@1) -> _pipe = increment_sync(Pl),
do_pipeline(_pipe, Ext, Rest, Sock@1) end
)
end
).
-file("src/pgl/internal/protocol.gleam", 638).
?DOC(false).
-spec handle_parameter_description(
pipeline(MMM),
list(pgl@internal@encode:'query'(MMM, any())),
extended(MMM),
list(integer()),
pgl@internal@socket:socket()
) -> {ok, pipeline(MMM)} | {error, pgl@internal:internal_error()}.
handle_parameter_description(Pl, Queries, Ext, Oids, Sock) ->
case Queries of
[] ->
do_pipeline(Pl, Ext, Queries, Sock);
[Query] ->
next_param_description(Pl, Query, [], Ext, Oids, Sock);
[Query@1 | Rest] ->
next_param_description(Pl, Query@1, Rest, Ext, Oids, Sock)
end.
-file("src/pgl/internal/protocol.gleam", 540).
?DOC(false).
-spec do_pipeline(
pipeline(MLX),
extended(MLX),
list(pgl@internal@encode:'query'(MLX, any())),
pgl@internal@socket:socket()
) -> {ok, pipeline(MLX)} | {error, pgl@internal:internal_error()}.
do_pipeline(Pl, Ext, Queries, Sock) ->
gleam@result:'try'(receive_message(Sock), fun(Msg) -> case Msg of
bind_complete ->
do_pipeline(Pl, Ext, Queries, Sock);
{command_complete, _, Count} ->
Ext@1 = {extended,
erlang:element(2, Ext),
erlang:element(3, Ext),
erlang:element(4, Ext),
erlang:element(5, Ext),
erlang:element(6, Ext),
erlang:element(7, Ext),
Count},
Acc = gleam@list:prepend(erlang:element(4, Pl), Ext@1),
Next_ext = {extended,
erlang:element(2, Ext@1),
erlang:element(3, Ext@1),
erlang:element(4, Ext@1),
[],
[],
[],
0},
_pipe = set_acc(Pl, Acc),
do_pipeline(_pipe, Next_ext, Queries, Sock);
{data_row, Values} ->
_pipe@1 = handle_data_row(
Values,
Ext,
erlang:element(3, Ext)
),
gleam@result:'try'(
_pipe@1,
fun(_capture) ->
do_pipeline(Pl, _capture, Queries, Sock)
end
);
{error_response, Fields} ->
_pipe@2 = Fields,
_pipe@3 = handle_error_response(_pipe@2),
error_response_cleanup(
_pipe@3,
erlang:element(2, Ext),
erlang:element(2, Pl),
erlang:element(3, Pl),
Sock
);
no_data ->
do_pipeline(Pl, Ext, Queries, Sock);
{notice_response, _} ->
do_pipeline(Pl, Ext, Queries, Sock);
{notification_response, _, _, _} ->
do_pipeline(Pl, Ext, Queries, Sock);
{parameter_description, _, Data_types} ->
handle_parameter_description(
Pl,
Queries,
Ext,
Data_types,
Sock
);
parse_complete ->
do_pipeline(Pl, Ext, Queries, Sock);
{ready_for_query, _} ->
Pl@1 = increment_ready(Pl),
case erlang:element(2, Pl@1) > erlang:element(3, Pl@1) of
true ->
do_pipeline(Pl@1, Ext, Queries, Sock);
false ->
{ok, reverse_acc(Pl@1)}
end;
{row_description, _, Descriptions} ->
_pipe@4 = handle_row_description(Ext, Descriptions),
do_pipeline(Pl, _pipe@4, Queries, Sock);
_ ->
_pipe@5 = sync(Sock),
_pipe@6 = gleam@result:try_recover(
_pipe@5,
fun(Field@0) -> {error, Field@0} end
),
gleam@result:'try'(
_pipe@6,
fun(_) ->
_pipe@7 = {protocol_error,
message_error,
<<"Unexpected message in flow"/utf8>>},
{error, _pipe@7}
end
)
end end).
-file("src/pgl/internal/protocol.gleam", 515).
?DOC(false).
-spec pipeline() -> pipeline(any()).
pipeline() ->
{pipeline, 0, 0, []}.
-file("src/pgl/internal/protocol.gleam", 420).
?DOC(false).
-spec process(
extended(MKD),
pgl@internal@encode:'query'(MKD, any()),
pgl@internal@socket:socket()
) -> {ok, extended(MKD)} | {error, pgl@internal:internal_error()}.
process(Flow, Query, Sock) ->
Needs_sync = pgl@internal@encode:needs_sync(Query),
gleam@result:'try'(
pgl@internal@encode:to_bit_array(Query),
fun(Packet) ->
Flow@1 = {extended,
Needs_sync,
erlang:element(3, Flow),
erlang:element(4, Flow),
erlang:element(5, Flow),
erlang:element(6, Flow),
erlang:element(7, Flow),
erlang:element(8, Flow)},
Pl = pipeline(),
gleam@result:'try'(
pgl@internal@socket:send(Sock, Packet),
fun(Sock@1) ->
gleam@result:'try'(
do_pipeline(Pl, Flow@1, [Query], Sock@1),
fun(Pl@1) -> case erlang:element(4, Pl@1) of
[Extended] ->
{ok, Extended};
_ ->
_pipe = {protocol_error,
processing_error,
<<"Missing rows"/utf8>>},
{error, _pipe}
end end
)
end
)
end
).
-file("src/pgl/internal/protocol.gleam", 519).
?DOC(false).
-spec batch_process(
pipeline(MLM),
extended(MLM),
list(pgl@internal@encode:'query'(MLM, any())),
pgl@internal@socket:socket()
) -> {ok, list(extended(MLM))} | {error, pgl@internal:internal_error()}.
batch_process(Flow, Extended, Queries, Sock) ->
gleam@result:'try'(
gleam@list:try_map(Queries, fun pgl@internal@encode:to_bit_array/1),
fun(Encoded) ->
Packet = begin
_pipe = Encoded,
_pipe@1 = gleam_stdlib:bit_array_concat(_pipe),
gleam@bit_array:append(_pipe@1, pgl@internal@encode:sync())
end,
gleam@result:'try'(
pgl@internal@socket:send(Sock, Packet),
fun(Sock@1) -> _pipe@2 = Flow,
_pipe@3 = increment_sync(_pipe@2),
_pipe@4 = do_pipeline(_pipe@3, Extended, Queries, Sock@1),
gleam@result:map(
_pipe@4,
fun(Pl) -> erlang:element(4, Pl) end
) end
)
end
).