// Helpers that turn backend message queues into higher-level client results.
///|
/// Construct a `RowStream` over one backend response queue.
fn new_row_stream(
background_group : @ref.Ref[@async.TaskGroup[Unit]?],
responses : @async.Queue[@backend.Message],
initial_columns? : Array[Column] = [],
cleanup? : StreamCleanup? = None,
) -> RowStream {
{
background_group,
columns: initial_columns,
responses,
cleanup,
command_tag: "",
row_count: 0,
suspended: false,
database_error: None,
detached: false,
summary: None,
}
}
///|
/// Construct a `SimpleQueryStream` over one backend response queue.
fn new_simple_query_stream(
background_group : @ref.Ref[@async.TaskGroup[Unit]?],
responses : @async.Queue[@backend.Message],
) -> SimpleQueryStream {
{
background_group,
responses,
columns: [],
database_error: None,
detached: false,
finished: false,
}
}
///|
/// Build the deferred cleanup object for a temporary prepared statement.
fn make_statement_cleanup(client : Client, name : Bytes) -> StreamCleanup raise {
{ client, bytes: close_bytes(b'S', name[:]), done: false, }
}
///|
/// Close a temporary prepared statement immediately.
///
/// This is used on early failures that happen before a `RowStream` can be
/// returned to the caller and therefore before deferred cleanup is available.
async fn close_temporary_statement(client : Client, name : Bytes) -> Unit {
let responses = client.send_request(Messages, close_bytes(b'S', name[:]))
drain_close_response(responses)
}
///|
/// Wait for PostgreSQL to acknowledge `COPY OUT` mode and build the stream.
async fn start_copy_out(
background_group : @ref.Ref[@async.TaskGroup[Unit]?],
responses : @async.Queue[@backend.Message],
) -> CopyOutStream {
let mut database_error : DatabaseError? = None
for res = responses.get() {
match res {
CopyOutResponse(body) =>
// The stream is valid only after PostgreSQL explicitly enters COPY OUT
// mode and advertises the per-column wire formats.
return {
background_group,
responses,
formats: parse_copy_formats(body.column_formats()),
database_error,
detached: false,
finished: false,
}
ErrorResponse(body) => {
database_error = Some(parse_database_error(body.fields()))
continue responses.get()
}
ReadyForQuery(_) =>
match database_error {
Some(err) => raise ClientError::Database(err)
None =>
raise ClientError::UnexpectedMessage("expected CopyOutResponse")
}
_ => raise ClientError::UnexpectedMessage("unexpected copy out response")
}
}
}
///|
/// Parse the server's response to a statement preparation request.
///
/// The server may send parameter descriptions, row descriptions, both, or
/// neither depending on the statement. Database errors are deferred until the
/// terminal `ReadyForQuery` just like the public stream helpers do.
async fn read_prepare_response(
shared : Shared,
responses : @async.Queue[@backend.Message],
) -> (Array[Type], Array[Column]) {
let mut param_oids : Array[@proto.Oid] = []
let mut columns : Array[Column] = []
let mut database_error : DatabaseError? = None
for res = responses.get() {
match res {
ParseComplete => continue responses.get()
ParameterDescription(body) => {
// Hold onto raw OIDs for now; the richer `Type` descriptors are loaded
// only after PostgreSQL finishes the prepare exchange successfully.
param_oids = parse_parameter_oids(body)
continue responses.get()
}
RowDescription(body) => {
columns = parse_columns(body)
continue responses.get()
}
NoData => continue responses.get()
ErrorResponse(body) => {
database_error = Some(parse_database_error(body.fields()))
continue responses.get()
}
ReadyForQuery(_) =>
match database_error {
Some(err) => raise ClientError::Database(err)
None => {
// Resolve cached type metadata once, at the end, so callers see a
// fully enriched result without partial-success bookkeeping.
let result = (
resolve_parameter_types(shared, param_oids),
resolve_columns(shared, columns),
)
break result
}
}
_ => raise ClientError::UnexpectedMessage("unexpected prepare response")
}
}
}
///|
/// Drain the backend completion sequence for COPY IN and return affected rows.
async fn collect_copy_in_completion(
responses : @async.Queue[@backend.Message],
) -> Int {
let mut command_tag = ""
let mut database_error : DatabaseError? = None
for res = responses.get() {
match res {
CommandComplete(body) => {
// COPY IN reports affected rows only through the trailing command tag,
// after all client input has already been sent.
command_tag = body.tag_str()
continue responses.get()
}
CopyDone => continue responses.get()
ErrorResponse(body) => {
database_error = Some(parse_database_error(body.fields()))
continue responses.get()
}
ReadyForQuery(_) =>
match database_error {
Some(err) => raise ClientError::Database(err)
None => break rows_affected(command_tag)
}
_ => raise ClientError::UnexpectedMessage("unexpected copy in completion")
}
}
}
///|
/// Drain the backend completion sequence for `Close`.
async fn drain_close_response(
responses : @async.Queue[@backend.Message],
) -> Unit {
let mut database_error : DatabaseError? = None
for res = responses.get() {
match res {
CloseComplete => continue responses.get()
ErrorResponse(body) => {
database_error = Some(parse_database_error(body.fields()))
continue responses.get()
}
ReadyForQuery(_) =>
match database_error {
Some(err) => raise ClientError::Database(err)
None => break ()
}
_ => raise ClientError::UnexpectedMessage("unexpected close response")
}
}
}
///|
/// Drain the response to a standalone `Sync`.
async fn drain_sync_response(
responses : @async.Queue[@backend.Message],
) -> Unit {
match responses.get() {
ReadyForQuery(_) => ()
_ => raise ClientError::UnexpectedMessage("unexpected sync response")
}
}
///|
/// Drain all responses produced by `batch_execute`.
async fn drain_simple_query(responses : @async.Queue[@backend.Message]) -> Unit {
let mut database_error : DatabaseError? = None
for res = responses.get() {
match res {
// `batch_execute` intentionally ignores per-statement rows and metadata;
// only terminal success or database failure matters here.
CommandComplete(_)
| EmptyQueryResponse
| RowDescription(_)
| DataRow(_) => continue responses.get()
ErrorResponse(body) => {
database_error = Some(parse_database_error(body.fields()))
continue responses.get()
}
ReadyForQuery(_) =>
match database_error {
Some(err) => raise ClientError::Database(err)
None => break ()
}
_ =>
raise ClientError::UnexpectedMessage(
"unexpected batch execute response",
)
}
}
}
///|
/// Extract the trailing numeric row count from a PostgreSQL command tag.
fn rows_affected(command_tag : String) -> Int {
match command_tag.split(" ").last() {
None => 0
Some(last) => @string.parse_int(last) catch { _ => 0 }
}
}
///|
/// Join transaction option clauses with spaces.
fn join_clauses(clauses : Array[String]) -> String {
let mut out = ""
for index, clause in clauses {
if index != 0 {
out += " "
}
out += clause
}
out
}