Current section

Files

Jump to
postgleam src postgleam@connection.erl
Raw

src/postgleam@connection.erl

-module(postgleam@connection).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch, inline]).
-define(FILEPATH, "src/postgleam/connection.gleam").
-export([send_message/2, send_bytes/2, receive_message/2, close_statement/3, sync_portal/2, prepare/5, disconnect/1, execute_portal/5, execute_prepared/5, extended_query/5, bind_and_execute_portal/6, simple_query/3, connect/1]).
-export_type([connection_state/0, simple_query_result/0, prepared_statement/0, extended_query_result/0, stream_chunk/0]).
-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 connection_state() :: {connection_state,
postgleam@internal@transport:transport(),
gleam@option:option(integer()),
gleam@option:option(integer()),
gleam@dict:dict(binary(), binary()),
postgleam@message:transaction_status(),
bitstring()}.
-type simple_query_result() :: {simple_query_result,
binary(),
list(binary()),
list(list(gleam@option:option(binary())))}.
-type prepared_statement() :: {prepared_statement,
binary(),
binary(),
list(integer()),
list(postgleam@message:row_field())}.
-type extended_query_result() :: {extended_query_result,
binary(),
list(postgleam@message:row_field()),
list(list(gleam@option:option(postgleam@value:value())))}.
-type stream_chunk() :: {stream_more,
list(list(gleam@option:option(postgleam@value:value())))} |
{stream_done,
binary(),
list(list(gleam@option:option(postgleam@value:value())))}.
-file("src/postgleam/connection.gleam", 82).
?DOC(" Send a frontend message over the transport\n").
-spec send_message(connection_state(), postgleam@message:frontend_message()) -> {ok,
connection_state()} |
{error, postgleam@error:error()}.
send_message(State, Msg) ->
Bytes = postgleam@message:encode_frontend(Msg),
case postgleam@internal@transport:send(erlang:element(2, State), Bytes) of
{ok, _} ->
{ok, State};
{error, E} ->
{error, E}
end.
-file("src/postgleam/connection.gleam", 94).
?DOC(" Send raw bytes over the transport\n").
-spec send_bytes(connection_state(), bitstring()) -> {ok, connection_state()} |
{error, postgleam@error:error()}.
send_bytes(State, Bytes) ->
case postgleam@internal@transport:send(erlang:element(2, State), Bytes) of
{ok, _} ->
{ok, State};
{error, E} ->
{error, E}
end.
-file("src/postgleam/connection.gleam", 112).
-spec receive_message_loop(connection_state(), integer()) -> {ok,
{postgleam@message:backend_message(), connection_state()}} |
{error, postgleam@error:error()}.
receive_message_loop(State, Timeout) ->
case postgleam@message:decode_backend(erlang:element(7, State)) of
{decoded, Msg, Rest} ->
{ok,
{Msg,
{connection_state,
erlang:element(2, State),
erlang:element(3, State),
erlang:element(4, State),
erlang:element(5, State),
erlang:element(6, State),
Rest}}};
{decode_failed, Reason} ->
{error, {protocol_error, <<"Decode failed: "/utf8, Reason/binary>>}};
incomplete ->
case postgleam@internal@transport:'receive'(
erlang:element(2, State),
Timeout
) of
{ok, Data} ->
New_buffer = <<(erlang:element(7, State))/bitstring,
Data/bitstring>>,
receive_message_loop(
{connection_state,
erlang:element(2, State),
erlang:element(3, State),
erlang:element(4, State),
erlang:element(5, State),
erlang:element(6, State),
New_buffer},
Timeout
);
{error, E} ->
{error, E}
end
end.
-file("src/postgleam/connection.gleam", 105).
?DOC(" Receive and decode the next backend message\n").
-spec receive_message(connection_state(), integer()) -> {ok,
{postgleam@message:backend_message(), connection_state()}} |
{error, postgleam@error:error()}.
receive_message(State, Timeout) ->
receive_message_loop(State, Timeout).
-file("src/postgleam/connection.gleam", 721).
-spec recv_close_complete(connection_state(), integer()) -> {ok,
connection_state()} |
{error, postgleam@error:error()}.
recv_close_complete(State, Timeout) ->
gleam@result:'try'(
receive_message(State, Timeout),
fun(_use0) ->
{Msg, State@1} = _use0,
case Msg of
close_complete ->
{ok, State@1};
{notice_response, _} ->
recv_close_complete(State@1, Timeout);
{parameter_status, Name, Val} ->
State@2 = {connection_state,
erlang:element(2, State@1),
erlang:element(3, State@1),
erlang:element(4, State@1),
gleam@dict:insert(erlang:element(5, State@1), Name, Val),
erlang:element(6, State@1),
erlang:element(7, State@1)},
recv_close_complete(State@2, Timeout);
_ ->
{error,
{protocol_error,
<<"Expected CloseComplete, got unexpected message"/utf8>>}}
end
end
).
-file("src/postgleam/connection.gleam", 744).
-spec recv_ready_for_query(connection_state(), binary(), integer()) -> {ok,
connection_state()} |
{error, postgleam@error:error()}.
recv_ready_for_query(State, Sql, Timeout) ->
gleam@result:'try'(
receive_message(State, Timeout),
fun(_use0) ->
{Msg, State@1} = _use0,
case Msg of
{ready_for_query, Status} ->
{ok,
{connection_state,
erlang:element(2, State@1),
erlang:element(3, State@1),
erlang:element(4, State@1),
erlang:element(5, State@1),
Status,
erlang:element(7, State@1)}};
{error_response, Fields} ->
Pg_fields = postgleam@error:parse_error_fields(Fields),
{error,
{pg_error,
Pg_fields,
erlang:element(3, State@1),
{some, Sql}}};
{notice_response, _} ->
recv_ready_for_query(State@1, Sql, Timeout);
{parameter_status, Name, Val} ->
State@2 = {connection_state,
erlang:element(2, State@1),
erlang:element(3, State@1),
erlang:element(4, State@1),
gleam@dict:insert(erlang:element(5, State@1), Name, Val),
erlang:element(6, State@1),
erlang:element(7, State@1)},
recv_ready_for_query(State@2, Sql, Timeout);
_ ->
{error,
{protocol_error,
<<"Expected ReadyForQuery, got unexpected message"/utf8>>}}
end
end
).
-file("src/postgleam/connection.gleam", 336).
?DOC(" Close a prepared statement: Close + Sync\n").
-spec close_statement(connection_state(), binary(), integer()) -> {ok,
connection_state()} |
{error, postgleam@error:error()}.
close_statement(State, Name, Timeout) ->
gleam@result:'try'(
send_message(State, {close, describe_statement, Name}),
fun(State@1) ->
gleam@result:'try'(
send_message(State@1, sync),
fun(State@2) ->
gleam@result:'try'(
recv_close_complete(State@2, Timeout),
fun(State@3) ->
recv_ready_for_query(State@3, <<""/utf8>>, Timeout)
end
)
end
)
end
).
-file("src/postgleam/connection.gleam", 492).
?DOC(" Finalize portal streaming: Sync to get ReadyForQuery\n").
-spec sync_portal(connection_state(), integer()) -> {ok, connection_state()} |
{error, postgleam@error:error()}.
sync_portal(State, Timeout) ->
gleam@result:'try'(
send_message(State, sync),
fun(State@1) -> recv_ready_for_query(State@1, <<""/utf8>>, Timeout) end
).
-file("src/postgleam/connection.gleam", 778).
?DOC(" Drain messages until ReadyForQuery (for error recovery)\n").
-spec drain_to_ready(connection_state(), integer()) -> {ok, connection_state()} |
{error, postgleam@error:error()}.
drain_to_ready(State, Timeout) ->
gleam@result:'try'(
receive_message(State, Timeout),
fun(_use0) ->
{Msg, State@1} = _use0,
case Msg of
{ready_for_query, Status} ->
{ok,
{connection_state,
erlang:element(2, State@1),
erlang:element(3, State@1),
erlang:element(4, State@1),
erlang:element(5, State@1),
Status,
erlang:element(7, State@1)}};
_ ->
drain_to_ready(State@1, Timeout)
end
end
).
-file("src/postgleam/connection.gleam", 502).
-spec recv_parse_complete(connection_state(), binary(), integer()) -> {ok,
connection_state()} |
{error, postgleam@error:error()}.
recv_parse_complete(State, Sql, Timeout) ->
gleam@result:'try'(
receive_message(State, Timeout),
fun(_use0) ->
{Msg, State@1} = _use0,
case Msg of
parse_complete ->
{ok, State@1};
{error_response, Fields} ->
Pg_fields = postgleam@error:parse_error_fields(Fields),
_ = drain_to_ready(State@1, Timeout),
{error,
{pg_error,
Pg_fields,
erlang:element(3, State@1),
{some, Sql}}};
{notice_response, _} ->
recv_parse_complete(State@1, Sql, Timeout);
{parameter_status, Name, Val} ->
State@2 = {connection_state,
erlang:element(2, State@1),
erlang:element(3, State@1),
erlang:element(4, State@1),
gleam@dict:insert(erlang:element(5, State@1), Name, Val),
erlang:element(6, State@1),
erlang:element(7, State@1)},
recv_parse_complete(State@2, Sql, Timeout);
_ ->
{error,
{protocol_error,
<<"Expected ParseComplete, got unexpected message"/utf8>>}}
end
end
).
-file("src/postgleam/connection.gleam", 536).
-spec recv_parameter_description(connection_state(), binary(), integer()) -> {ok,
{list(integer()), connection_state()}} |
{error, postgleam@error:error()}.
recv_parameter_description(State, Sql, Timeout) ->
gleam@result:'try'(
receive_message(State, Timeout),
fun(_use0) ->
{Msg, State@1} = _use0,
case Msg of
{parameter_description, Type_oids} ->
{ok, {Type_oids, State@1}};
{error_response, Fields} ->
Pg_fields = postgleam@error:parse_error_fields(Fields),
_ = drain_to_ready(State@1, Timeout),
{error,
{pg_error,
Pg_fields,
erlang:element(3, State@1),
{some, Sql}}};
{notice_response, _} ->
recv_parameter_description(State@1, Sql, Timeout);
{parameter_status, Name, Val} ->
State@2 = {connection_state,
erlang:element(2, State@1),
erlang:element(3, State@1),
erlang:element(4, State@1),
gleam@dict:insert(erlang:element(5, State@1), Name, Val),
erlang:element(6, State@1),
erlang:element(7, State@1)},
recv_parameter_description(State@2, Sql, Timeout);
_ ->
{error,
{protocol_error,
<<"Expected ParameterDescription, got unexpected message"/utf8>>}}
end
end
).
-file("src/postgleam/connection.gleam", 569).
-spec recv_row_description_or_nodata(connection_state(), binary(), integer()) -> {ok,
{list(postgleam@message:row_field()), connection_state()}} |
{error, postgleam@error:error()}.
recv_row_description_or_nodata(State, Sql, Timeout) ->
gleam@result:'try'(
receive_message(State, Timeout),
fun(_use0) ->
{Msg, State@1} = _use0,
case Msg of
{row_description, Fields} ->
{ok, {Fields, State@1}};
no_data ->
{ok, {[], State@1}};
{error_response, Fields@1} ->
Pg_fields = postgleam@error:parse_error_fields(Fields@1),
_ = drain_to_ready(State@1, Timeout),
{error,
{pg_error,
Pg_fields,
erlang:element(3, State@1),
{some, Sql}}};
{notice_response, _} ->
recv_row_description_or_nodata(State@1, Sql, Timeout);
{parameter_status, Name, Val} ->
State@2 = {connection_state,
erlang:element(2, State@1),
erlang:element(3, State@1),
erlang:element(4, State@1),
gleam@dict:insert(erlang:element(5, State@1), Name, Val),
erlang:element(6, State@1),
erlang:element(7, State@1)},
recv_row_description_or_nodata(State@2, Sql, Timeout);
_ ->
{error,
{protocol_error,
<<"Expected RowDescription or NoData, got unexpected message"/utf8>>}}
end
end
).
-file("src/postgleam/connection.gleam", 238).
?DOC(
" Prepare a statement: Parse + Describe + Sync\n"
" Returns the prepared statement with parameter and result type info.\n"
).
-spec prepare(
connection_state(),
binary(),
binary(),
list(integer()),
integer()
) -> {ok, {prepared_statement(), connection_state()}} |
{error, postgleam@error:error()}.
prepare(State, Name, Sql, Type_oids, Timeout) ->
gleam@result:'try'(
send_message(State, {parse, Name, Sql, Type_oids}),
fun(State@1) ->
gleam@result:'try'(
send_message(State@1, {describe, describe_statement, Name}),
fun(State@2) ->
gleam@result:'try'(
send_message(State@2, sync),
fun(State@3) ->
gleam@result:'try'(
recv_parse_complete(State@3, Sql, Timeout),
fun(State@4) ->
gleam@result:'try'(
recv_parameter_description(
State@4,
Sql,
Timeout
),
fun(_use0) ->
{Param_oids_result, State@5} = _use0,
gleam@result:'try'(
recv_row_description_or_nodata(
State@5,
Sql,
Timeout
),
fun(_use0@1) ->
{Result_fields, State@6} = _use0@1,
gleam@result:'try'(
recv_ready_for_query(
State@6,
Sql,
Timeout
),
fun(State@7) ->
Prepared = {prepared_statement,
Name,
Sql,
Param_oids_result,
Result_fields},
{ok,
{Prepared,
State@7}}
end
)
end
)
end
)
end
)
end
)
end
)
end
).
-file("src/postgleam/connection.gleam", 604).
-spec recv_bind_complete(connection_state(), binary(), integer()) -> {ok,
connection_state()} |
{error, postgleam@error:error()}.
recv_bind_complete(State, Sql, Timeout) ->
gleam@result:'try'(
receive_message(State, Timeout),
fun(_use0) ->
{Msg, State@1} = _use0,
case Msg of
bind_complete ->
{ok, State@1};
{error_response, Fields} ->
Pg_fields = postgleam@error:parse_error_fields(Fields),
_ = drain_to_ready(State@1, Timeout),
{error,
{pg_error,
Pg_fields,
erlang:element(3, State@1),
{some, Sql}}};
{notice_response, _} ->
recv_bind_complete(State@1, Sql, Timeout);
{parameter_status, Name, Val} ->
State@2 = {connection_state,
erlang:element(2, State@1),
erlang:element(3, State@1),
erlang:element(4, State@1),
gleam@dict:insert(erlang:element(5, State@1), Name, Val),
erlang:element(6, State@1),
erlang:element(7, State@1)},
recv_bind_complete(State@2, Sql, Timeout);
_ ->
{error,
{protocol_error,
<<"Expected BindComplete, got unexpected message"/utf8>>}}
end
end
).
-file("src/postgleam/connection.gleam", 835).
?DOC(" Resolve codecs for result fields once, to avoid per-row registry lookups\n").
-spec resolve_row_decoders(
list(postgleam@message:row_field()),
gleam@dict:dict(integer(), postgleam@codec:codec())
) -> list(fun((bitstring()) -> {ok, postgleam@value:value()} | {error, binary()})).
resolve_row_decoders(Fields, Reg) ->
gleam@list:map(
Fields,
fun(Field) ->
case postgleam@codec@registry:lookup(Reg, erlang:element(5, Field)) of
{ok, Codec} ->
erlang:element(6, Codec);
{error, _} ->
fun(Bytes) -> case gleam@bit_array:to_string(Bytes) of
{ok, S} ->
{ok, {text, S}};
{error, _} ->
{ok, {bytea, Bytes}}
end end
end
end
).
-file("src/postgleam/connection.gleam", 861).
-spec decode_values_with_decoders(
list(gleam@option:option(bitstring())),
list(fun((bitstring()) -> {ok, postgleam@value:value()} | {error, binary()}))
) -> list(gleam@option:option(postgleam@value:value())).
decode_values_with_decoders(Raw_values, Decoders) ->
case {Raw_values, Decoders} of
{[], []} ->
[];
{[none | Rest_vals], [_ | Rest_decoders]} ->
[none | decode_values_with_decoders(Rest_vals, Rest_decoders)];
{[{some, Bytes} | Rest_vals@1], [Decoder | Rest_decoders@1]} ->
Decoded = case Decoder(Bytes) of
{ok, Val} ->
{some, Val};
{error, _} ->
{some, {text, <<"<decode error>"/utf8>>}}
end,
[Decoded |
decode_values_with_decoders(Rest_vals@1, Rest_decoders@1)];
{_, _} ->
[]
end.
-file("src/postgleam/connection.gleam", 853).
?DOC(" Decode a binary DataRow using pre-resolved decoders\n").
-spec decode_binary_row(
bitstring(),
list(fun((bitstring()) -> {ok, postgleam@value:value()} | {error, binary()}))
) -> list(gleam@option:option(postgleam@value:value())).
decode_binary_row(Values, Decoders) ->
Raw_values = postgleam@message:extract_row_values(Values),
decode_values_with_decoders(Raw_values, Decoders).
-file("src/postgleam/connection.gleam", 637).
-spec recv_execute_rows(
connection_state(),
prepared_statement(),
list(fun((bitstring()) -> {ok, postgleam@value:value()} | {error, binary()})),
integer(),
list(list(gleam@option:option(postgleam@value:value())))
) -> {ok,
{binary(),
list(list(gleam@option:option(postgleam@value:value()))),
connection_state()}} |
{error, postgleam@error:error()}.
recv_execute_rows(State, Prepared, Decoders, Timeout, Rows) ->
gleam@result:'try'(
receive_message(State, Timeout),
fun(_use0) ->
{Msg, State@1} = _use0,
case Msg of
{data_row, Values} ->
Row = decode_binary_row(Values, Decoders),
recv_execute_rows(
State@1,
Prepared,
Decoders,
Timeout,
[Row | Rows]
);
{command_complete, Tag} ->
{ok, {Tag, Rows, State@1}};
empty_query_response ->
{ok, {<<""/utf8>>, Rows, State@1}};
{error_response, Fields} ->
Pg_fields = postgleam@error:parse_error_fields(Fields),
_ = drain_to_ready(State@1, Timeout),
{error,
{pg_error,
Pg_fields,
erlang:element(3, State@1),
{some, erlang:element(3, Prepared)}}};
{notice_response, _} ->
recv_execute_rows(
State@1,
Prepared,
Decoders,
Timeout,
Rows
);
{parameter_status, Name, Val} ->
State@2 = {connection_state,
erlang:element(2, State@1),
erlang:element(3, State@1),
erlang:element(4, State@1),
gleam@dict:insert(erlang:element(5, State@1), Name, Val),
erlang:element(6, State@1),
erlang:element(7, State@1)},
recv_execute_rows(
State@2,
Prepared,
Decoders,
Timeout,
Rows
);
_ ->
{error,
{protocol_error,
<<"Expected DataRow or CommandComplete, got unexpected message"/utf8>>}}
end
end
).
-file("src/postgleam/connection.gleam", 976).
?DOC(" Process startup messages after authentication until ReadyForQuery\n").
-spec startup_loop(connection_state(), postgleam@config:config()) -> {ok,
connection_state()} |
{error, postgleam@error:error()}.
startup_loop(State, Config) ->
gleam@result:'try'(
receive_message(State, erlang:element(7, Config)),
fun(_use0) ->
{Msg, State@1} = _use0,
case Msg of
{ready_for_query, Status} ->
{ok,
{connection_state,
erlang:element(2, State@1),
erlang:element(3, State@1),
erlang:element(4, State@1),
erlang:element(5, State@1),
Status,
erlang:element(7, State@1)}};
{parameter_status, Name, Value} ->
State@2 = {connection_state,
erlang:element(2, State@1),
erlang:element(3, State@1),
erlang:element(4, State@1),
gleam@dict:insert(
erlang:element(5, State@1),
Name,
Value
),
erlang:element(6, State@1),
erlang:element(7, State@1)},
startup_loop(State@2, Config);
{backend_key_data, Pid, Key} ->
State@3 = {connection_state,
erlang:element(2, State@1),
{some, Pid},
{some, Key},
erlang:element(5, State@1),
erlang:element(6, State@1),
erlang:element(7, State@1)},
startup_loop(State@3, Config);
{notice_response, _} ->
startup_loop(State@1, Config);
{error_response, Fields} ->
Pg_fields = postgleam@error:parse_error_fields(Fields),
{error,
{pg_error, Pg_fields, erlang:element(3, State@1), none}};
_ ->
{error,
{protocol_error,
<<"Unexpected message during startup"/utf8>>}}
end
end
).
-file("src/postgleam/connection.gleam", 1015).
?DOC(" Close the connection\n").
-spec disconnect(connection_state()) -> nil.
disconnect(State) ->
_ = send_message(State, terminate),
postgleam@internal@transport:close(erlang:element(2, State)).
-file("src/postgleam/connection.gleam", 1022).
-spec extract_column_names(list(postgleam@message:row_field())) -> list(binary()).
extract_column_names(Fields) ->
case Fields of
[] ->
[];
[F | Rest] ->
[erlang:element(2, F) | extract_column_names(Rest)]
end.
-file("src/postgleam/connection.gleam", 1044).
-spec connect_error_to_string(mug:connect_error()) -> binary().
connect_error_to_string(Err) ->
case Err of
{connect_failed_ipv4, _} ->
<<"connection failed (IPv4)"/utf8>>;
{connect_failed_ipv6, _} ->
<<"connection failed (IPv6)"/utf8>>;
{connect_failed_both, _, _} ->
<<"connection failed (IPv4 and IPv6)"/utf8>>
end.
-file("src/postgleam/connection.gleam", 1056).
-spec list_reverse_loop(list(GNZ), list(GNZ)) -> list(GNZ).
list_reverse_loop(L, Acc) ->
case L of
[] ->
Acc;
[X | Rest] ->
list_reverse_loop(Rest, [X | Acc])
end.
-file("src/postgleam/connection.gleam", 1052).
-spec list_reverse(list(GNW)) -> list(GNW).
list_reverse(L) ->
list_reverse_loop(L, []).
-file("src/postgleam/connection.gleam", 678).
-spec recv_stream_rows(
connection_state(),
prepared_statement(),
list(fun((bitstring()) -> {ok, postgleam@value:value()} | {error, binary()})),
integer(),
list(list(gleam@option:option(postgleam@value:value())))
) -> {ok, {stream_chunk(), connection_state()}} |
{error, postgleam@error:error()}.
recv_stream_rows(State, Prepared, Decoders, Timeout, Rows) ->
gleam@result:'try'(
receive_message(State, Timeout),
fun(_use0) ->
{Msg, State@1} = _use0,
case Msg of
{data_row, Values} ->
Row = decode_binary_row(Values, Decoders),
recv_stream_rows(
State@1,
Prepared,
Decoders,
Timeout,
[Row | Rows]
);
{command_complete, Tag} ->
{ok, {{stream_done, Tag, list_reverse(Rows)}, State@1}};
portal_suspended ->
{ok, {{stream_more, list_reverse(Rows)}, State@1}};
empty_query_response ->
{ok,
{{stream_done, <<""/utf8>>, list_reverse(Rows)},
State@1}};
{error_response, Fields} ->
Pg_fields = postgleam@error:parse_error_fields(Fields),
_ = drain_to_ready(State@1, Timeout),
{error,
{pg_error,
Pg_fields,
erlang:element(3, State@1),
{some, erlang:element(3, Prepared)}}};
{notice_response, _} ->
recv_stream_rows(State@1, Prepared, Decoders, Timeout, Rows);
{parameter_status, Name, Val} ->
State@2 = {connection_state,
erlang:element(2, State@1),
erlang:element(3, State@1),
erlang:element(4, State@1),
gleam@dict:insert(erlang:element(5, State@1), Name, Val),
erlang:element(6, State@1),
erlang:element(7, State@1)},
recv_stream_rows(State@2, Prepared, Decoders, Timeout, Rows);
_ ->
{error,
{protocol_error,
<<"Expected DataRow, CommandComplete, or PortalSuspended"/utf8>>}}
end
end
).
-file("src/postgleam/connection.gleam", 475).
?DOC(
" Execute the next chunk from a portal: Execute(max_rows) + Flush\n"
" Returns StreamMore if portal is suspended (more rows), StreamDone if complete.\n"
).
-spec execute_portal(
connection_state(),
prepared_statement(),
gleam@dict:dict(integer(), postgleam@codec:codec()),
integer(),
integer()
) -> {ok, {stream_chunk(), connection_state()}} |
{error, postgleam@error:error()}.
execute_portal(State, Prepared, Registry, Max_rows, Timeout) ->
gleam@result:'try'(
send_message(State, {execute, <<""/utf8>>, Max_rows}),
fun(State@1) ->
gleam@result:'try'(
send_message(State@1, flush),
fun(State@2) ->
Decoders = resolve_row_decoders(
erlang:element(5, Prepared),
Registry
),
recv_stream_rows(State@2, Prepared, Decoders, Timeout, [])
end
)
end
).
-file("src/postgleam/connection.gleam", 801).
-spec encode_params_loop(
list(gleam@option:option(postgleam@value:value())),
list(integer()),
gleam@dict:dict(integer(), postgleam@codec:codec()),
list(gleam@option:option(bitstring()))
) -> {ok, list(gleam@option:option(bitstring()))} |
{error, postgleam@error:error()}.
encode_params_loop(Params, Oids, Reg, Acc) ->
case {Params, Oids} of
{[], []} ->
{ok, list_reverse(Acc)};
{[none | Rest_params], [_ | Rest_oids]} ->
encode_params_loop(Rest_params, Rest_oids, Reg, [none | Acc]);
{[{some, Val} | Rest_params@1], [Oid | Rest_oids@1]} ->
case postgleam@codec@registry:lookup(Reg, Oid) of
{ok, Codec} ->
case (erlang:element(5, Codec))(Val) of
{ok, Bytes} ->
encode_params_loop(
Rest_params@1,
Rest_oids@1,
Reg,
[{some, Bytes} | Acc]
);
{error, E} ->
{error, {encode_error, E}}
end;
{error, E@1} ->
{error, {encode_error, E@1}}
end;
{_, _} ->
{error,
{encode_error,
<<"Parameter count mismatch: params and OIDs have different lengths"/utf8>>}}
end.
-file("src/postgleam/connection.gleam", 793).
?DOC(" Encode parameters to binary using the codec registry\n").
-spec encode_params_binary(
list(gleam@option:option(postgleam@value:value())),
list(integer()),
gleam@dict:dict(integer(), postgleam@codec:codec())
) -> {ok, list(gleam@option:option(bitstring()))} |
{error, postgleam@error:error()}.
encode_params_binary(Params, Param_oids, Reg) ->
encode_params_loop(Params, Param_oids, Reg, []).
-file("src/postgleam/connection.gleam", 281).
?DOC(
" Execute a prepared statement with binary parameters: Bind + Execute + Sync\n"
" Parameters are encoded using the codec registry.\n"
).
-spec execute_prepared(
connection_state(),
prepared_statement(),
list(gleam@option:option(postgleam@value:value())),
gleam@dict:dict(integer(), postgleam@codec:codec()),
integer()
) -> {ok, {extended_query_result(), connection_state()}} |
{error, postgleam@error:error()}.
execute_prepared(State, Prepared, Params, Registry, Timeout) ->
gleam@result:'try'(
encode_params_binary(Params, erlang:element(4, Prepared), Registry),
fun(Encoded_params) ->
Param_formats = gleam@list:map(
erlang:element(4, Prepared),
fun(_) -> binary_format end
),
Result_formats = gleam@list:map(
erlang:element(5, Prepared),
fun(_) -> binary_format end
),
gleam@result:'try'(
send_message(
State,
{bind,
<<""/utf8>>,
erlang:element(2, Prepared),
Param_formats,
Encoded_params,
Result_formats}
),
fun(State@1) ->
gleam@result:'try'(
send_message(State@1, {execute, <<""/utf8>>, 0}),
fun(State@2) ->
gleam@result:'try'(
send_message(State@2, sync),
fun(State@3) ->
gleam@result:'try'(
recv_bind_complete(
State@3,
erlang:element(3, Prepared),
Timeout
),
fun(State@4) ->
Decoders = resolve_row_decoders(
erlang:element(5, Prepared),
Registry
),
gleam@result:'try'(
recv_execute_rows(
State@4,
Prepared,
Decoders,
Timeout,
[]
),
fun(_use0) ->
{Tag, Rows, State@5} = _use0,
gleam@result:'try'(
recv_ready_for_query(
State@5,
erlang:element(
3,
Prepared
),
Timeout
),
fun(State@6) ->
{ok,
{{extended_query_result,
Tag,
erlang:element(
5,
Prepared
),
list_reverse(
Rows
)},
State@6}}
end
)
end
)
end
)
end
)
end
)
end
)
end
).
-file("src/postgleam/connection.gleam", 355).
?DOC(
" All-in-one extended query: Parse + Describe + Bind + Execute + Close + Sync\n"
" Uses unnamed statement and portal for simplicity.\n"
).
-spec extended_query(
connection_state(),
binary(),
list(gleam@option:option(postgleam@value:value())),
gleam@dict:dict(integer(), postgleam@codec:codec()),
integer()
) -> {ok, {extended_query_result(), connection_state()}} |
{error, postgleam@error:error()}.
extended_query(State, Sql, Params, Registry, Timeout) ->
gleam@result:'try'(
send_message(State, {parse, <<""/utf8>>, Sql, []}),
fun(State@1) ->
gleam@result:'try'(
send_message(
State@1,
{describe, describe_statement, <<""/utf8>>}
),
fun(State@2) ->
gleam@result:'try'(
send_message(State@2, sync),
fun(State@3) ->
gleam@result:'try'(
recv_parse_complete(State@3, Sql, Timeout),
fun(State@4) ->
gleam@result:'try'(
recv_parameter_description(
State@4,
Sql,
Timeout
),
fun(_use0) ->
{Param_oids, State@5} = _use0,
gleam@result:'try'(
recv_row_description_or_nodata(
State@5,
Sql,
Timeout
),
fun(_use0@1) ->
{Result_fields, State@6} = _use0@1,
gleam@result:'try'(
recv_ready_for_query(
State@6,
Sql,
Timeout
),
fun(State@7) ->
gleam@result:'try'(
encode_params_binary(
Params,
Param_oids,
Registry
),
fun(
Encoded_params
) ->
Param_formats = gleam@list:map(
Param_oids,
fun(_) ->
binary_format
end
),
Result_formats = gleam@list:map(
Result_fields,
fun(_) ->
binary_format
end
),
Prepared = {prepared_statement,
<<""/utf8>>,
Sql,
Param_oids,
Result_fields},
gleam@result:'try'(
send_message(
State@7,
{bind,
<<""/utf8>>,
<<""/utf8>>,
Param_formats,
Encoded_params,
Result_formats}
),
fun(
State@8
) ->
gleam@result:'try'(
send_message(
State@8,
{execute,
<<""/utf8>>,
0}
),
fun(
State@9
) ->
gleam@result:'try'(
send_message(
State@9,
{close,
describe_statement,
<<""/utf8>>}
),
fun(
State@10
) ->
gleam@result:'try'(
send_message(
State@10,
sync
),
fun(
State@11
) ->
gleam@result:'try'(
recv_bind_complete(
State@11,
Sql,
Timeout
),
fun(
State@12
) ->
Decoders = resolve_row_decoders(
erlang:element(
5,
Prepared
),
Registry
),
gleam@result:'try'(
recv_execute_rows(
State@12,
Prepared,
Decoders,
Timeout,
[]
),
fun(
_use0@2
) ->
{Tag,
Rows,
State@13} = _use0@2,
gleam@result:'try'(
recv_close_complete(
State@13,
Timeout
),
fun(
State@14
) ->
gleam@result:'try'(
recv_ready_for_query(
State@14,
Sql,
Timeout
),
fun(
State@15
) ->
{ok,
{{extended_query_result,
Tag,
Result_fields,
list_reverse(
Rows
)},
State@15}}
end
)
end
)
end
)
end
)
end
)
end
)
end
)
end
)
end
)
end
)
end
)
end
)
end
)
end
)
end
)
end
).
-file("src/postgleam/connection.gleam", 441).
?DOC(
" Bind parameters and execute first chunk: Bind + Execute(max_rows) + Flush\n"
" Uses Flush instead of Sync to keep the unnamed portal alive.\n"
" Returns the first chunk of rows.\n"
).
-spec bind_and_execute_portal(
connection_state(),
prepared_statement(),
list(gleam@option:option(postgleam@value:value())),
gleam@dict:dict(integer(), postgleam@codec:codec()),
integer(),
integer()
) -> {ok, {stream_chunk(), connection_state()}} |
{error, postgleam@error:error()}.
bind_and_execute_portal(State, Prepared, Params, Registry, Max_rows, Timeout) ->
gleam@result:'try'(
encode_params_binary(Params, erlang:element(4, Prepared), Registry),
fun(Encoded_params) ->
Param_formats = gleam@list:map(
erlang:element(4, Prepared),
fun(_) -> binary_format end
),
Result_formats = gleam@list:map(
erlang:element(5, Prepared),
fun(_) -> binary_format end
),
gleam@result:'try'(
send_message(
State,
{bind,
<<""/utf8>>,
erlang:element(2, Prepared),
Param_formats,
Encoded_params,
Result_formats}
),
fun(State@1) ->
gleam@result:'try'(
send_message(State@1, {execute, <<""/utf8>>, Max_rows}),
fun(State@2) ->
gleam@result:'try'(
send_message(State@2, flush),
fun(State@3) ->
gleam@result:'try'(
recv_bind_complete(
State@3,
erlang:element(3, Prepared),
Timeout
),
fun(State@4) ->
Decoders = resolve_row_decoders(
erlang:element(5, Prepared),
Registry
),
recv_stream_rows(
State@4,
Prepared,
Decoders,
Timeout,
[]
)
end
)
end
)
end
)
end
)
end
).
-file("src/postgleam/connection.gleam", 1063).
-spec list_map(list(GOD), fun((GOD) -> GOF)) -> list(GOF).
list_map(L, F) ->
case L of
[] ->
[];
[X | Rest] ->
[F(X) | list_map(Rest, F)]
end.
-file("src/postgleam/connection.gleam", 1029).
-spec decode_text_row(bitstring()) -> list(gleam@option:option(binary())).
decode_text_row(Values) ->
_pipe = postgleam@message:extract_row_values(Values),
list_map(_pipe, fun(V) -> case V of
{some, Bytes} ->
case gleam@bit_array:to_string(Bytes) of
{ok, S} ->
{some, S};
{error, _} ->
{some, <<"<binary>"/utf8>>}
end;
none ->
none
end end).
-file("src/postgleam/connection.gleam", 157).
-spec simple_query_loop(
connection_state(),
binary(),
integer(),
list(simple_query_result()),
gleam@option:option(list(binary())),
list(list(gleam@option:option(binary())))
) -> {ok, {list(simple_query_result()), connection_state()}} |
{error, postgleam@error:error()}.
simple_query_loop(State, Sql, Timeout, Results, Current_columns, Current_rows) ->
gleam@result:'try'(
receive_message(State, Timeout),
fun(_use0) ->
{Msg, State@1} = _use0,
case Msg of
{ready_for_query, Status} ->
State@2 = {connection_state,
erlang:element(2, State@1),
erlang:element(3, State@1),
erlang:element(4, State@1),
erlang:element(5, State@1),
Status,
erlang:element(7, State@1)},
{ok, {list_reverse(Results), State@2}};
{row_description, Fields} ->
Columns = extract_column_names(Fields),
simple_query_loop(
State@1,
Sql,
Timeout,
Results,
{some, Columns},
[]
);
{data_row, Values} ->
Row = decode_text_row(Values),
simple_query_loop(
State@1,
Sql,
Timeout,
Results,
Current_columns,
[Row | Current_rows]
);
{command_complete, Tag} ->
Cols = case Current_columns of
{some, C} ->
C;
none ->
[]
end,
Result = {simple_query_result,
Tag,
Cols,
list_reverse(Current_rows)},
simple_query_loop(
State@1,
Sql,
Timeout,
[Result | Results],
none,
[]
);
empty_query_response ->
simple_query_loop(
State@1,
Sql,
Timeout,
Results,
Current_columns,
Current_rows
);
{error_response, Fields@1} ->
Pg_fields = postgleam@error:parse_error_fields(Fields@1),
_ = drain_to_ready(State@1, Timeout),
{error,
{pg_error,
Pg_fields,
erlang:element(3, State@1),
{some, Sql}}};
{notice_response, _} ->
simple_query_loop(
State@1,
Sql,
Timeout,
Results,
Current_columns,
Current_rows
);
{notification_response, _, _, _} ->
simple_query_loop(
State@1,
Sql,
Timeout,
Results,
Current_columns,
Current_rows
);
_ ->
{error,
{protocol_error,
<<"Unexpected message during simple query"/utf8>>}}
end
end
).
-file("src/postgleam/connection.gleam", 139).
?DOC(" Execute a simple query (text protocol) and return all results\n").
-spec simple_query(connection_state(), binary(), integer()) -> {ok,
{list(simple_query_result()), connection_state()}} |
{error, postgleam@error:error()}.
simple_query(State, Sql, Timeout) ->
gleam@result:'try'(
send_message(State, {simple_query, Sql}),
fun(State@1) ->
simple_query_loop(State@1, Sql, Timeout, [], none, [])
end
).
-file("src/postgleam/connection.gleam", 951).
-spec scram_final_loop(
connection_state(),
postgleam@config:config(),
postgleam@auth@scram:scram_state()
) -> {ok, connection_state()} | {error, postgleam@error:error()}.
scram_final_loop(State, Config, Scram_state) ->
gleam@result:'try'(
receive_message(State, erlang:element(7, Config)),
fun(_use0) ->
{Msg, State@1} = _use0,
case Msg of
{authentication_msg, {auth_sasl_final, Server_final}} ->
case postgleam@auth@scram:verify_server(
Server_final,
Scram_state
) of
{ok, _} ->
authenticate_loop(State@1, Config);
{error, E} ->
{error, {authentication_error, E}}
end;
{error_response, Fields} ->
Pg_fields = postgleam@error:parse_error_fields(Fields),
{error,
{authentication_error, erlang:element(4, Pg_fields)}};
_ ->
{error,
{protocol_error,
<<"Expected SASL final, got unexpected message"/utf8>>}}
end
end
).
-file("src/postgleam/connection.gleam", 882).
-spec authenticate_loop(connection_state(), postgleam@config:config()) -> {ok,
connection_state()} |
{error, postgleam@error:error()}.
authenticate_loop(State, Config) ->
gleam@result:'try'(
receive_message(State, erlang:element(7, Config)),
fun(_use0) ->
{Msg, State@1} = _use0,
case Msg of
{authentication_msg, auth_ok} ->
startup_loop(State@1, Config);
{authentication_msg, auth_cleartext} ->
gleam@result:'try'(
send_message(
State@1,
{password_message, erlang:element(6, Config)}
),
fun(State@2) -> authenticate_loop(State@2, Config) end
);
{authentication_msg, {auth_md5, Salt}} ->
Hash = postgleam@auth@md5:hash_password(
erlang:element(6, Config),
erlang:element(5, Config),
Salt
),
gleam@result:'try'(
send_message(State@1, {password_message, Hash}),
fun(State@3) -> authenticate_loop(State@3, Config) end
);
{authentication_msg, {auth_sasl, _}} ->
{Mechanism, Client_first_data} = postgleam@auth@scram:client_first(
),
Client_first_bare@1 = case postgleam@auth@scram:extract_client_first_bare(
Client_first_data
) of
{ok, Client_first_bare} -> Client_first_bare;
_assert_fail ->
erlang:error(#{gleam_error => let_assert,
message => <<"Pattern match failed, no pattern matched the value."/utf8>>,
file => <<?FILEPATH/utf8>>,
module => <<"postgleam/connection"/utf8>>,
function => <<"authenticate_loop"/utf8>>,
line => 902,
value => _assert_fail,
start => 27654,
'end' => 27747,
pattern_start => 27665,
pattern_end => 27686})
end,
gleam@result:'try'(
send_message(
State@1,
{s_a_s_l_initial_response,
Mechanism,
Client_first_data}
),
fun(State@4) ->
scram_continue_loop(
State@4,
Config,
Client_first_bare@1
)
end
);
{error_response, Fields} ->
Pg_fields = postgleam@error:parse_error_fields(Fields),
{error,
{authentication_error,
<<"Authentication failed: "/utf8,
(erlang:element(4, Pg_fields))/binary>>}};
_ ->
{error,
{protocol_error,
<<"Unexpected message during authentication"/utf8>>}}
end
end
).
-file("src/postgleam/connection.gleam", 922).
-spec scram_continue_loop(
connection_state(),
postgleam@config:config(),
binary()
) -> {ok, connection_state()} | {error, postgleam@error:error()}.
scram_continue_loop(State, Config, Client_first_bare) ->
gleam@result:'try'(
receive_message(State, erlang:element(7, Config)),
fun(_use0) ->
{Msg, State@1} = _use0,
case Msg of
{authentication_msg, {auth_sasl_continue, Server_first}} ->
case postgleam@auth@scram:client_final(
Server_first,
Client_first_bare,
erlang:element(6, Config)
) of
{ok, {Client_final, Scram_state}} ->
gleam@result:'try'(
send_message(
State@1,
{s_a_s_l_response, Client_final}
),
fun(State@2) ->
scram_final_loop(
State@2,
Config,
Scram_state
)
end
);
{error, E} ->
{error,
{authentication_error,
<<"SCRAM error: "/utf8, E/binary>>}}
end;
{error_response, Fields} ->
Pg_fields = postgleam@error:parse_error_fields(Fields),
{error,
{authentication_error, erlang:element(4, Pg_fields)}};
_ ->
{error,
{protocol_error,
<<"Expected SASL continue, got unexpected message"/utf8>>}}
end
end
).
-file("src/postgleam/connection.gleam", 38).
?DOC(" Establish a connection to PostgreSQL, complete authentication, and return ready state\n").
-spec connect(postgleam@config:config()) -> {ok, connection_state()} |
{error, postgleam@error:error()}.
connect(Config) ->
gleam@result:'try'(
begin
_pipe = mug:new(
erlang:element(2, Config),
erlang:element(3, Config)
),
_pipe@1 = mug:timeout(_pipe, erlang:element(8, Config)),
_pipe@2 = mug:connect(_pipe@1),
gleam@result:map_error(
_pipe@2,
fun(E) ->
{socket_error,
<<"TCP connect failed: "/utf8,
(connect_error_to_string(E))/binary>>}
end
)
end,
fun(Tcp_socket) -> gleam@result:'try'(case erlang:element(10, Config) of
ssl_disabled ->
{ok, {tcp, Tcp_socket}};
ssl_verified ->
postgleam@internal@transport:upgrade_to_ssl(
Tcp_socket,
erlang:element(2, Config),
erlang:element(8, Config),
true
);
ssl_unverified ->
postgleam@internal@transport:upgrade_to_ssl(
Tcp_socket,
erlang:element(2, Config),
erlang:element(8, Config),
false
)
end, fun(Conn_transport) ->
State = {connection_state,
Conn_transport,
none,
none,
maps:new(),
idle,
<<>>},
Base_params = [{<<"user"/utf8>>, erlang:element(5, Config)},
{<<"database"/utf8>>, erlang:element(4, Config)}],
Startup = {startup_message,
lists:append(Base_params, erlang:element(9, Config))},
gleam@result:'try'(
send_message(State, Startup),
fun(State@1) -> authenticate_loop(State@1, Config) end
)
end) end
).