Skip to content
Merged
Show file tree
Hide file tree
Changes from 2 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
129 changes: 55 additions & 74 deletions src/sql_jsc/postgres/PostgresSQLQuery.rs
Original file line number Diff line number Diff line change
Expand Up @@ -479,15 +479,6 @@ impl PostgresSQLQuery {
};
let connection: &PostgresSQLConnection = &connection;

// `KeepAlive::ref_` takes an `EventLoopCtx` (manual vtable in `bun_io`), not a
// `*mut VirtualMachine`. `global_object.bun_vm()` and `get_vm_ctx(.Js)` both
// resolve to the same singleton JS VM, so route through the global hook —
// identical to `PostgresSQLConnection::vm_ctx`.
connection.poll_ref.with_mut(|r| {
r.ref_(bun_io::posix_event_loop::get_vm_ctx(
bun_io::AllocatorType::Js,
))
});
let query = arguments[1];

if !query.is_object() {
Expand All @@ -501,6 +492,26 @@ impl PostgresSQLQuery {
let writer = connection.writer();
// We need a strong reference to the query so that it doesn't get GC'd
this.ref_();
// Shared cleanup for every error-return path below: drop any statement
// ref this query took plus the speculative `ref_()` above.
let release_query_ref = || {
this.release_statement();
// SAFETY: undoes the speculative `this.ref_()` above; count was ≥2, never frees here.
unsafe { Self::deref(this_ptr) };
};
// Shared error tail: throw `err` as a postgres error unless an exception
// is already pending.
let throw_write_error = |msg: &[u8], err: AnyPostgresError| -> JsError {
if !global_object.has_exception() {
return global_object.throw_value(postgres_error_to_js(
global_object,
Some(msg),
err,
));
}
JsError::Thrown
};

if this.flags.get().simple {
bun_core::scoped_log!(Postgres, "executeQuery");

Expand All @@ -520,20 +531,8 @@ impl PostgresSQLQuery {
let can_execute = !connection.has_query_running();
if can_execute {
if let Err(err) = PostgresRequest::execute_query(query_str.slice(), writer) {
// fail to run do cleanup — sole owner just created above
// (rc=1); `release_statement` decrements → 0 frees.
this.release_statement();
// SAFETY: undoes the speculative `this.ref_()` above; count was ≥2, never frees here.
unsafe { Self::deref(this_ptr) };

if !global_object.has_exception() {
return Err(global_object.throw_value(postgres_error_to_js(
global_object,
Some(b"failed to execute query"),
err,
)));
}
return Err(JsError::Thrown);
release_query_ref();
return Err(throw_write_error(b"failed to execute query", err));
}
{
let mut f = connection.flags.get();
Expand All @@ -552,15 +551,20 @@ impl PostgresSQLQuery {
.with_mut(|q| q.write_item(this_ptr))
.is_err()
{
// fail to run do cleanup — sole owner just created above
// (rc=1); `release_statement` decrements → 0 frees.
this.release_statement();
// SAFETY: undoes the speculative `this.ref_()` above; count was ≥2, never frees here.
unsafe { Self::deref(this_ptr) };

release_query_ref();
return Err(global_object.throw_out_of_memory());
}

// Request is enqueued: keep the event loop alive until the server
// responds. KeepAlive is a flag (not a count), so taking this any
// earlier would leave it stuck Active on the synchronous-error
// returns above.
connection.poll_ref.with_mut(|r| {
r.ref_(bun_io::posix_event_loop::get_vm_ctx(
bun_io::AllocatorType::Js,
))
});

this.this_value.with_mut(|r| r.upgrade(global_object));
js::target_set_cached(this_value, global_object, query);
if this.status.get() == Status::Running {
Expand Down Expand Up @@ -630,6 +634,7 @@ impl PostgresSQLQuery {
Ok(v) => v,
Err(err) => {
drop(signature);
release_query_ref();
return Err(
global_object.throw_error(err.into(), "failed to allocate statement")
);
Expand Down Expand Up @@ -670,21 +675,11 @@ impl PostgresSQLQuery {
columns_value,
writer,
) {
// fail to run do cleanup — drop the ref we took above.
this.release_statement();
// SAFETY: undoes the speculative `this.ref_()` above; count was ≥2, never frees here.
unsafe { Self::deref(this_ptr) };

if !global_object.has_exception() {
return Err(global_object.throw_value(
postgres_error_to_js(
global_object,
Some(b"failed to bind and execute query"),
err,
),
));
}
return Err(JsError::Thrown);
release_query_ref();
return Err(throw_write_error(
b"failed to bind and execute query",
err,
));
}
{
let mut f = connection.flags.get();
Expand Down Expand Up @@ -726,17 +721,8 @@ impl PostgresSQLQuery {
.with_mut(|m| m.remove(&signature_hash));
}
drop(signature);
this.release_statement();
// SAFETY: undoes the speculative `this.ref_()` above; count was ≥2, never frees here.
unsafe { Self::deref(this_ptr) };
if !global_object.has_exception() {
return Err(global_object.throw_value(postgres_error_to_js(
global_object,
Some(b"failed to prepare and query"),
err,
)));
}
return Err(JsError::Thrown);
release_query_ref();
return Err(throw_write_error(b"failed to prepare and query", err));
}
{
let mut f = connection.flags.get();
Expand Down Expand Up @@ -767,17 +753,8 @@ impl PostgresSQLQuery {
.with_mut(|m| m.remove(&signature_hash));
}
drop(signature);
this.release_statement();
// SAFETY: undoes the speculative `this.ref_()` above; count was ≥2, never frees here.
unsafe { Self::deref(this_ptr) };
if !global_object.has_exception() {
return Err(global_object.throw_value(postgres_error_to_js(
global_object,
Some(b"failed to write query"),
err,
)));
}
return Err(JsError::Thrown);
release_query_ref();
return Err(throw_write_error(b"failed to write query", err));
}
if let Err(err) = writer.write(&protocol::SYNC) {
if connection_entry_value.is_some() {
Expand All @@ -786,14 +763,8 @@ impl PostgresSQLQuery {
.with_mut(|m| m.remove(&signature_hash));
}
drop(signature);
if !global_object.has_exception() {
return Err(global_object.throw_value(postgres_error_to_js(
global_object,
Some(b"failed to flush"),
err,
)));
}
return Err(JsError::Thrown);
release_query_ref();
return Err(throw_write_error(b"failed to flush", err));
}
{
let mut f = connection.flags.get();
Expand Down Expand Up @@ -854,8 +825,18 @@ impl PostgresSQLQuery {
.with_mut(|q| q.write_item(this_ptr))
.is_err()
{
release_query_ref();
return Err(global_object.throw_out_of_memory());
}
// Request is enqueued: keep the event loop alive until the server
// responds. See the matching call in the simple-query branch above
// for why this must come after every fallible step.
connection.poll_ref.with_mut(|r| {
r.ref_(bun_io::posix_event_loop::get_vm_ctx(
bun_io::AllocatorType::Js,
))
});

this.this_value.with_mut(|r| r.upgrade(global_object));

js::target_set_cached(this_value, global_object, query);
Expand Down
58 changes: 58 additions & 0 deletions test/js/sql/sql-postgres-run-error-pollref-fixture.ts

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

32 changes: 32 additions & 0 deletions test/js/sql/sql-postgres-run-error-pollref.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,32 @@
// PostgresSQLQuery.do_run refs the connection's poll_ref KeepAlive. KeepAlive
// is a two-state flag, not a counter, so when this query is the only in-flight
// work the call flips Inactive -> Active. When do_run then returns early with
// a synchronous error (bad binding, signature-generation failure, OOM during
// enqueue, ...) the poll_ref must not be left Active: nothing else on the
// connection will touch it until the next server message, so the event loop
// stays pinned and the process never exits.
//
// The fixture connects to a mock server, lets the connection go idle, then
// issues a query whose binding is rejected synchronously before anything is
// written. It must print the rejection and exit on its own.

import { expect, test } from "bun:test";
import { bunEnv, bunExe } from "harness";
import path from "node:path";

test("postgres: synchronous do_run failure does not pin the event loop", async () => {
await using proc = Bun.spawn({
cmd: [bunExe(), path.join(import.meta.dir, "sql-postgres-run-error-pollref-fixture.ts")],
env: bunEnv,
stdout: "pipe",
stderr: "pipe",
});

const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]);
void stderr;

expect(stdout).toBe("rejected:ERR_INVALID_ARG_TYPE\n");
// exited on its own, not killed by the runner's timeout
expect(proc.signalCode).toBeNull();
expect(exitCode).toBe(0);
});
Loading