diff --git a/src/sql/shared/ConnectionFlags.rs b/src/sql/shared/ConnectionFlags.rs index b06ab88756ae..6e6946a42514 100644 --- a/src/sql/shared/ConnectionFlags.rs +++ b/src/sql/shared/ConnectionFlags.rs @@ -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`. + const IS_DISPATCHING = 1 << 5; } } diff --git a/src/sql_jsc/postgres/PostgresSQLConnection.rs b/src/sql_jsc/postgres/PostgresSQLConnection.rs index d740a1421cdb..e388eff01d1b 100644 --- a/src/sql_jsc/postgres/PostgresSQLConnection.rs +++ b/src/sql_jsc/postgres/PostgresSQLConnection.rs @@ -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() { @@ -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. + 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. + pub(crate) fn while_dispatching(&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 @@ -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. 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; diff --git a/src/sql_jsc/postgres/PostgresSQLQuery.rs b/src/sql_jsc/postgres/PostgresSQLQuery.rs index 54e1e2d1954c..76fa23839d52 100644 --- a/src/sql_jsc/postgres/PostgresSQLQuery.rs +++ b/src/sql_jsc/postgres/PostgresSQLQuery.rs @@ -522,6 +522,9 @@ impl PostgresSQLQuery { } JsError::Thrown }; + // Dispatched from inside another request's encoder: enqueue only + // (PostgresSQLConnection::while_dispatching). + let dispatching = connection.is_dispatching(); if this.flags.get().simple { bun_core::scoped_log!(Postgres, "executeQuery"); @@ -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(); @@ -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 @@ -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. + 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 @@ -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 @@ -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()); diff --git a/test/js/sql/postgres-dispatch-during-bind-fixture.ts b/test/js/sql/postgres-dispatch-during-bind-fixture.ts new file mode 100644 index 000000000000..32929c733d7a --- /dev/null +++ b/test/js/sql/postgres-dispatch-during-bind-fixture.ts @@ -0,0 +1,195 @@ +// Fixture for postgres-dispatch-during-bind.test.ts. Runs one scenario (named +// by SCENARIO) against the server at DATABASE_URL in its own process, because +// before the fix several of these scenarios abort the process (panic in the +// Bind encoder) or wedge the connection instead of failing a single assertion. +// The test kills a wedged fixture when it times out. +// +// Every scenario binds a parameter whose conversion (valueOf(), or toString() +// where the parameter goes out in text format) dispatches more queries on the +// same (max: 1) connection while the outer query's Bind message is being +// encoded, and reports what each query settled with plus how many times the +// parameter was converted. One conversion per query is the observable side of +// "one Bind message per request". +import { SQL } from "bun"; + +const sql = new SQL({ url: process.env.DATABASE_URL!, max: 1, idleTimeout: 30 }); +await sql.connect(); + +// The int4 parameters below are all plain objects, so within a scenario the +// warm-up query and the outer query map to the same prepared statement (the +// statement name encodes the parameters' JS types). `::int4` makes the server +// type the parameter as int4, which is what gets valueOf() called while the +// Bind is encoded. +const int = (n: number) => ({ valueOf: () => n }); + +let conversions = 0; +const dispatched: Promise[] = []; +/** int4 parameter whose first conversion runs `dispatch` (synchronously, inside the Bind encoder). */ +function dispatching(value: number, dispatch: () => void) { + let fired = false; + return { + valueOf() { + conversions++; + if (!fired) { + fired = true; + dispatch(); + } + return value; + }, + }; +} + +function settle(promise: Promise) { + return promise.then( + ok => ({ ok }), + (e: any) => ({ err: e?.code ?? e?.message ?? String(e) }), + ); +} + +async function report(outer: Promise) { + return { + outer: await settle(outer), + dispatched: await Promise.all(dispatched.map(settle)), + conversions, + }; +} + +const scenarios: Record Promise> = { + // First execution of the statement: Parse round trip first, then the Bind is + // written from advance() in the ReadyForQuery handler. The nested query is a + // new statement text, so it is enqueued and tries to drain the queue itself. + async "first execution, nested new statement"() { + return report(sql`select ${dispatching(1, () => dispatched.push(sql`select 2 as y`.execute()))}::int4 as x`); + }, + + // Same outer shape; the nested query reuses a statement prepared earlier, the + // shape that writes its Bind at enqueue time. + async "first execution, nested prepared statement"() { + await sql`select ${"warm"}::text as t`; + return report( + sql`select ${dispatching(1, () => dispatched.push(sql`select ${"nested"}::text as t`.execute()))}::int4 as x`, + ); + }, + + // Outer statement already prepared: its Bind is written at enqueue time, from + // run(), with the request not yet in the queue. The nested query is a new + // statement text. + async "prepared statement, nested new statement"() { + await sql`select ${int(0)}::int4 as x`; + return report(sql`select ${dispatching(1, () => dispatched.push(sql`select 2 as y`.execute()))}::int4 as x`); + }, + + // Both outer and nested take the enqueue-time path; the nested one is the + // very same statement. + async "prepared statement, nested same statement"() { + await sql`select ${int(0)}::int4 as x`; + return report( + sql`select ${dispatching(1, () => dispatched.push(sql`select ${int(2)}::int4 as x`.execute()))}::int4 as x`, + ); + }, + + // One conversion dispatches a burst: prepared statements, a new statement + // text and a simple-protocol query. All of them have to come back in order. + async "prepared statement, nested burst"() { + await sql`select ${int(0)}::int4 as x`; + await sql`select ${"warm"}::text as t`; + return report( + sql`select ${dispatching(1, () => { + dispatched.push(sql`select ${"a"}::text as t`.execute()); + dispatched.push(sql`select ${int(2)}::int4 as x`.execute()); + dispatched.push(sql`select 3 as y`.execute()); + dispatched.push(sql.unsafe("select 'simple' as s").execute()); + dispatched.push(sql`select ${"b"}::text as t`.execute()); + })}::int4 as x`, + ); + }, + + // The nested query's own parameter dispatches yet another query when its + // Bind is encoded in turn. + async "nested query dispatches again from its own bind"() { + const third = dispatching(3, () => {}); + const second = dispatching(2, () => dispatched.push(sql`select ${third}::int4 as x`.execute())); + return report( + sql`select ${dispatching(1, () => dispatched.push(sql`select ${second}::int4 as x`.execute()))}::int4 as x`, + ); + }, + + // Inside a transaction the nested query goes straight to the reserved + // connection. + async "inside a transaction"() { + let nested: Promise | undefined; + const rows = await settle( + sql.begin(async tx => { + const outer = await tx`select ${dispatching(1, () => (nested = tx`select 2 as y`.execute()))}::int4 as x`; + return [outer, await nested]; + }), + ); + return { rows, conversions }; + }, + + // prepare: false sends Parse+Bind+Execute in one batch from advance(); the + // parameter is sent in text format there, so the conversion hook is + // toString() rather than valueOf(). + async "unnamed statements (prepare: false)"() { + await using unprepared = new SQL({ url: process.env.DATABASE_URL!, max: 1, idleTimeout: 30, prepare: false }); + const param = { + toString() { + conversions++; + if (dispatched.length === 0) { + dispatched.push(unprepared`select 2 as y`.execute()); + dispatched.push(unprepared`select ${"nested"}::text as t`.execute()); + } + return "1"; + }, + }; + // Awaited here so that `unprepared` is not disposed before the queries settle. + const result = await report(unprepared`select ${param}::int4 as x`); + return result; + }, + + // The conversion dispatches a query and then throws. The outer query rejects + // with that error; the dispatched one was only enqueued and must still be + // dispatched afterwards rather than sit in the queue forever. Only its + // settling is asserted: as long as an aborted Bind leaves its torn message in + // the write buffer the server drops the connection and it rejects; once that + // is rolled back it resolves. + async "prepared statement, conversion throws after dispatching"() { + await sql`select ${int(0)}::int4 as x`; + return throwAfterDispatching(); + }, + + // Same, with another query already in flight on the connection, so the outer + // request is not at the head of the queue when its Bind fails and has to be + // swept out from behind the in-flight one. + async "prepared statement behind an in-flight query, conversion throws after dispatching"() { + await sql`select ${int(0)}::int4 as x`; + const ahead = sql`select ${int(7)}::int4 as x`.execute(); + // Issues the outer query synchronously, while `ahead` is still in flight. + const rest = throwAfterDispatching(); + return { ahead: await settle(ahead), ...(await rest) }; + }, +}; + +/** Issues the outer query synchronously; everything after the first await is reporting. */ +async function throwAfterDispatching() { + const param = { + valueOf() { + conversions++; + dispatched.push(sql`select 2 as y`.execute()); + throw new RangeError("boom"); + }, + }; + const outer = await settle(sql`select ${param}::int4 as x`); + const settled = await Promise.all(dispatched.map(settle)); + // The pool reconnects if the torn message cost it the connection. + const afterwards = await settle(sql`select 3 as z`); + return { outer, dispatchedSettled: settled.length, afterwards, conversions }; +} + +const scenario = scenarios[process.env.SCENARIO!]; +if (!scenario) { + console.log(JSON.stringify({ error: `unknown scenario ${process.env.SCENARIO}` })); + process.exit(1); +} +console.log(JSON.stringify(await scenario())); +await sql.close(); diff --git a/test/js/sql/postgres-dispatch-during-bind.test.ts b/test/js/sql/postgres-dispatch-during-bind.test.ts new file mode 100644 index 000000000000..93a3c4a9e642 --- /dev/null +++ b/test/js/sql/postgres-dispatch-during-bind.test.ts @@ -0,0 +1,121 @@ +// Encoding a Bind message converts each parameter through user JS (valueOf / +// toString / toJSON), and that JS can synchronously dispatch another query on +// the same connection: with max: 1 (or inside a transaction) the pool hands it +// the connection whose write buffer currently ends in the half-encoded Bind. +// +// The nested dispatch used to either write its own messages into the middle of +// that Bind, or re-enter advance(), which bound the still-Pending outer request +// a second time and then flushed the buffer out from under the outer encoder, +// whose length prefix was an offset into it: +// +// panic: range start index 42 out of range for slice of length 4 +// (Writer::pwrite, PostgresSQLConnection.rs; debug builds trip +// "pending_requests underflow" first) +// +// For an already-prepared statement the Bind is written at enqueue time instead, +// and the nested query's frames and queue entry both ended up ahead of the +// outer's half-written Bind, so the server's error for the torn message was +// delivered to the nested query and the outer one never settled. +// +// Now a dispatch that lands inside an encoder only enqueues, and whoever is +// encoding drains the queue afterwards. Each scenario runs in a subprocess and +// reports what every query settled with and how many times the outer parameter +// was converted (one conversion per query is the observable form of one Bind +// per request); a broken scenario shows up as a panic in the subprocess's +// stderr or, for the wedged-connection shapes, as the test timing out. +import { expect, test } from "bun:test"; +import { bunEnv, bunExe, describeWithContainer } from "harness"; +import path from "node:path"; + +const fixture = path.join(import.meta.dir, "postgres-dispatch-during-bind-fixture.ts"); + +const scenarios: Record = { + "first execution, nested new statement": { + outer: { ok: [{ x: 1 }] }, + dispatched: [{ ok: [{ y: 2 }] }], + conversions: 1, + }, + "first execution, nested prepared statement": { + outer: { ok: [{ x: 1 }] }, + dispatched: [{ ok: [{ t: "nested" }] }], + conversions: 1, + }, + "prepared statement, nested new statement": { + outer: { ok: [{ x: 1 }] }, + dispatched: [{ ok: [{ y: 2 }] }], + conversions: 1, + }, + "prepared statement, nested same statement": { + outer: { ok: [{ x: 1 }] }, + dispatched: [{ ok: [{ x: 2 }] }], + conversions: 1, + }, + "prepared statement, nested burst": { + outer: { ok: [{ x: 1 }] }, + dispatched: [ + { ok: [{ t: "a" }] }, + { ok: [{ x: 2 }] }, + { ok: [{ y: 3 }] }, + { ok: [{ s: "simple" }] }, + { ok: [{ t: "b" }] }, + ], + conversions: 1, + }, + "nested query dispatches again from its own bind": { + outer: { ok: [{ x: 1 }] }, + dispatched: [{ ok: [{ x: 2 }] }, { ok: [{ x: 3 }] }], + conversions: 3, + }, + "inside a transaction": { + rows: { ok: [[{ x: 1 }], [{ y: 2 }]] }, + conversions: 1, + }, + "unnamed statements (prepare: false)": { + outer: { ok: [{ x: 1 }] }, + dispatched: [{ ok: [{ y: 2 }] }, { ok: [{ t: "nested" }] }], + conversions: 1, + }, + "prepared statement, conversion throws after dispatching": { + outer: { err: "boom" }, + dispatchedSettled: 1, + afterwards: { ok: [{ z: 3 }] }, + conversions: 1, + }, + "prepared statement behind an in-flight query, conversion throws after dispatching": { + ahead: { ok: [{ x: 7 }] }, + outer: { err: "boom" }, + dispatchedSettled: 1, + afterwards: { ok: [{ z: 3 }] }, + conversions: 1, + }, +}; + +describeWithContainer("postgres", { image: "postgres_plain" }, container => { + for (const [scenario, expected] of Object.entries(scenarios)) { + test.concurrent(scenario, async () => { + await container.ready; + await using proc = Bun.spawn({ + cmd: [bunExe(), fixture], + env: { + ...bunEnv, + DATABASE_URL: `postgres://bun_sql_test@${container.host}:${container.port}/bun_sql_test`, + SCENARIO: scenario, + }, + stdout: "pipe", + stderr: "pipe", + }); + const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); + + let result: unknown = stdout; + try { + result = JSON.parse(stdout); + } catch {} + expect({ result, exitCode, stderr }).toEqual({ + result: expected, + exitCode: 0, + // not asserted; included so a panic trace shows up in the diff + stderr: expect.any(String), + }); + }); + } +});