Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions src/sql/shared/ConnectionFlags.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,9 @@ bitflags! {
const USE_UNNAMED_PREPARED_STATEMENTS = 1 << 2;
const WAITING_TO_PREPARE = 1 << 3;
const HAS_BACKPRESSURE = 1 << 4;
/// Requests are being encoded into the write buffer; see
/// `PostgresSQLConnection::while_dispatching`.
Comment thread
robobun marked this conversation as resolved.
const IS_DISPATCHING = 1 << 5;
}
}

Expand Down
58 changes: 57 additions & 1 deletion src/sql_jsc/postgres/PostgresSQLConnection.rs
Original file line number Diff line number Diff line change
Expand Up @@ -661,11 +661,17 @@ impl PostgresSQLConnection {
}

pub(crate) fn flush_data(&self) {
let flags = self.flags.get();
// we know we still have backpressure so just return we will flush later
if self.flags.get().contains(ConnectionFlags::HAS_BACKPRESSURE) {
if flags.contains(ConnectionFlags::HAS_BACKPRESSURE) {
debug!("flushData: has backpressure");
return;
}
// the buffer may end in a half-encoded message; the encoder flushes when it is done
if flags.contains(ConnectionFlags::IS_DISPATCHING) {
debug!("flushData: dispatching");
return;
}

let chunk = self.write_buffer.get().remaining();
if chunk.is_empty() {
Expand Down Expand Up @@ -1555,6 +1561,42 @@ impl PostgresSQLConnection {
self.requests.with_mut(|q| q.discard(1));
}

/// Same disposal `advance()` gives a request it failed to encode: popped if
/// it heads the queue, otherwise (marked `Fail`) swept by a later `advance()`.
/// Then dispatches whatever was enqueued during the encode, so the caller
/// must have taken the encoder's exception off the VM.
Comment thread
robobun marked this conversation as resolved.
pub(crate) fn discard_failed_request(&self, request: *mut PostgresSQLQuery) {
debug_assert!(!self.is_dispatching());
self.discard_request(request);
if self.pending_requests.get() > 0 {
self.advance_and_flush();
}
}

#[inline]
pub(crate) fn is_dispatching(&self) -> bool {
self.flags.get().contains(ConnectionFlags::IS_DISPATCHING)
}

/// Runs `encode` (appends requests' messages to `write_buffer`) with
/// [`ConnectionFlags::IS_DISPATCHING`] set.
///
/// Encoding a Bind converts each parameter through user JS, which can
/// synchronously dispatch another query on this connection. While the flag
/// is set such a dispatch only enqueues (`PostgresSQLQuery::do_run`),
/// `advance()` returns at once and `flush_data()` is a no-op: the buffer may
/// end in a message whose length prefix has yet to be patched, and the
/// request in progress still looks unwritten to `advance()`. The caller
/// drains and flushes once `encode` returns.
Comment thread
robobun marked this conversation as resolved.
pub(crate) fn while_dispatching<T>(&self, encode: impl FnOnce() -> T) -> T {
debug_assert!(!self.is_dispatching());
self.update_flags(|f| f.insert(ConnectionFlags::IS_DISPATCHING));
scopeguard::defer! {
self.update_flags(|f| f.remove(ConnectionFlags::IS_DISPATCHING));
}
encode()
}

pub(crate) fn has_query_running(&self) -> bool {
!self
.flags
Expand Down Expand Up @@ -1808,7 +1850,21 @@ impl PostgresSQLConnection {
}
}

/// Writes out as many queued requests as the connection state allows.
///
/// A call made from inside an encoder ([`Self::while_dispatching`]) returns
/// at once; the newly enqueued request is reached by the outer loop (it
/// re-reads the queue length every iteration) or by the ReadyForQuery that
/// ends the request being encoded.
Comment thread
robobun marked this conversation as resolved.
fn advance(&self) {
if self.is_dispatching() {
debug!("advance: already dispatching");
return;
}
self.while_dispatching(|| self.advance_impl());
}

fn advance_impl(&self) {
let mut offset: usize = 0;
debug!("advance");
// The cleanup loop runs after the main loop returns;
Expand Down
68 changes: 48 additions & 20 deletions src/sql_jsc/postgres/PostgresSQLQuery.rs
Original file line number Diff line number Diff line change
Expand Up @@ -522,6 +522,9 @@ impl PostgresSQLQuery {
}
JsError::Thrown
};
// Dispatched from inside another request's encoder: enqueue only
// (PostgresSQLConnection::while_dispatching).
Comment thread
robobun marked this conversation as resolved.
let dispatching = connection.is_dispatching();

if this.flags.get().simple {
bun_core::scoped_log!(Postgres, "executeQuery");
Expand All @@ -539,7 +542,7 @@ impl PostgresSQLQuery {
// Query is simple and it's the only owner of the statement
this.statement.set(Some(stmt));

let can_execute = !connection.has_query_running();
let can_execute = !dispatching && !connection.has_query_running();
if can_execute {
if let Err(err) = PostgresRequest::execute_query(query_str.slice(), writer) {
release_query_ref();
Expand Down Expand Up @@ -617,6 +620,7 @@ impl PostgresSQLQuery {

let has_params = signature.fields.len() > 0;
let mut did_write = false;
let mut enqueued = false;
'enqueue: {
// Note: `connection_entry_value` is a *mut into connection.statements value slot;
// holding a `&mut` across other &mut connection borrows below trips borrowck, so
Expand Down Expand Up @@ -661,32 +665,55 @@ impl PostgresSQLQuery {
// request has already emitted its bytes; otherwise this
// Bind+Execute would overtake an earlier unwritten
// request on the wire while reply attribution stays FIFO.
if (!connection.has_query_running() || connection.can_pipeline())
if !dispatching
&& (!connection.has_query_running() || connection.can_pipeline())
&& connection.pending_requests.get() == 0
{
this.update_flags(|f| f.binary = !stmt.fields.is_empty());
bun_core::scoped_log!(Postgres, "bindAndExecute");

// bindAndExecute will bind + execute, it will change to running after binding is complete
if let Err(err) = PostgresRequest::bind_and_execute(
global_object,
stmt,
binding_value,
columns_value,
writer,
) {
// Enqueued before encoding so that queries dispatched
// from user JS during the encode queue up behind it
// (replies are attributed in queue order); Binding
// because it is not counted in pending_requests.
Comment thread
robobun marked this conversation as resolved.
if connection
.requests
.with_mut(|q| q.write_item(this_ptr))
.is_err()
{
release_query_ref();
return Err(throw_write_error(
b"failed to bind and execute query",
err,
));
return Err(global_object.throw_out_of_memory());
}
enqueued = true;
this.status.set(Status::Binding);

// bindAndExecute will bind + execute, it will change to running after binding is complete
if let Err(err) = connection.while_dispatching(|| {
PostgresRequest::bind_and_execute(
global_object,
stmt,
binding_value,
columns_value,
writer,
)
}) {
// taken while discard_failed_request dispatches, rethrown below
let exception = global_object.try_take_exception();
this.status.set(Status::Fail);
connection.discard_failed_request(this_ptr);
return Err(match exception {
Some(exception) => global_object.throw_value(exception),
None => throw_write_error(
b"failed to bind and execute query",
err,
),
});
}
{
let mut f = connection.flags.get();
f.set(ConnectionFlags::IS_READY_FOR_QUERY, false);
connection.flags.set(f);
}
this.status.set(Status::Binding);
this.update_flags(|f| f.counter = RequestCounter::Pipelined);
connection
.pipelined_requests
Expand Down Expand Up @@ -719,7 +746,7 @@ impl PostgresSQLQuery {
};
connection_entry_value = Some(entry_value_ptr);
}
let can_execute = !connection.has_query_running();
let can_execute = !dispatching && !connection.has_query_running();

if can_execute {
// If it does not have params, we can write and execute immediately in one go
Expand Down Expand Up @@ -838,10 +865,11 @@ impl PostgresSQLQuery {
}
}

if connection
.requests
.with_mut(|q| q.write_item(this_ptr))
.is_err()
if !enqueued
&& connection
.requests
.with_mut(|q| q.write_item(this_ptr))
.is_err()
{
release_query_ref();
return Err(global_object.throw_out_of_memory());
Expand Down
Loading
Loading