diff --git a/docs/runtime/sql.mdx b/docs/runtime/sql.mdx index bfe28bc9fea6..1ec843def489 100644 --- a/docs/runtime/sql.mdx +++ b/docs/runtime/sql.mdx @@ -985,6 +985,8 @@ const sql = new SQL({ }); ``` +Each connection keeps up to 256 idle prepared statements cached; a statement still referenced by a query (in-flight or completed but not yet garbage-collected) is never evicted, so the count can temporarily exceed that. Past the limit, the least recently used idle statement is deallocated on the server with a `Close` message to make room. Without the cap, a long-lived connection running many distinct query strings (for example, an ORM interpolating identifiers) would accumulate named prepared statements in the server session, and their metadata in the client, until the connection closed. + When `prepare: false` is set: Queries still use the "extended" protocol, but run as [unnamed prepared statements](https://www.postgresql.org/docs/current/protocol-flow.html#PROTOCOL-FLOW-EXT-QUERY). An unnamed prepared statement lasts only until the next Parse statement specifying the unnamed statement as destination is issued. diff --git a/src/sql_jsc/postgres/PostgresSQLConnection.rs b/src/sql_jsc/postgres/PostgresSQLConnection.rs index f42d7c179971..0950cb723bc8 100644 --- a/src/sql_jsc/postgres/PostgresSQLConnection.rs +++ b/src/sql_jsc/postgres/PostgresSQLConnection.rs @@ -53,6 +53,12 @@ bun_core::define_scoped_log!(debug, Postgres, visible); const MAX_PIPELINE_SIZE: usize = u16::MAX as usize; // about 64KB per connection +/// Per-connection cap on cached named prepared statements. Each cached +/// statement pins its metadata and row `Structure` on the client and a named +/// prepared statement in the server session until it is closed, so the cache +/// must be bounded and evictions must send `Close`. +const MAX_CACHED_PREPARED_STATEMENTS: usize = 256; + type PreparedStatementsMap = StringHashMap<*mut PostgresSQLStatement>; pub mod js { @@ -125,6 +131,9 @@ pub struct PostgresSQLConnection { // so `vm_mut()`'s `&mut *as_ptr()` is sound. pub vm: BackRef, pub statements: JsCell, + /// Monotonic clock for the statement cache's LRU order; bumped on every + /// cache lookup and stamped into `PostgresSQLStatement::last_used` on reuse. + statement_lru_clock: Cell, pub prepared_statement_id: Cell, pub pending_activity_count: AtomicU32, // Self-wrapper back-ref (the JS object that owns this payload). Stored as a @@ -1206,6 +1215,7 @@ pub(crate) fn call(global_object: &JSGlobalObject, callframe: &CallFrame) -> JsR // canonical `VirtualMachine::as_mut()` accessor. vm: BackRef::new_mut(vm), statements: JsCell::new(PreparedStatementsMap::default()), + statement_lru_clock: Cell::new(0), prepared_statement_id: Cell::new(0), pending_activity_count: AtomicU32::new(0), js_value: JsCell::new(crate::jsc::JsRef::empty()), @@ -1634,6 +1644,110 @@ impl PostgresSQLConnection { && !flags.contains(ConnectionFlags::HAS_BACKPRESSURE) // dont make sense to buffer more if we have backpressure && (self.write_buffer.get().len() as usize) < MAX_PIPELINE_SIZE // buffer is too big need to flush before pipeline more } + + /// Statement-cache hit probe, keyed by the query signature. Stamps the LRU + /// clock on the hit so the least recently *reused* entry is the eviction + /// victim. Returns `None` on a miss; the caller then goes through + /// [`put_statement`](Self::put_statement). + pub(crate) fn lookup_statement(&self, name: &[u8]) -> Option<*mut PostgresSQLStatement> { + let stmt_ptr = self.statements.get().get(name).copied()?; + self.statement_lru_clock + .set(self.statement_lru_clock.get() + 1); + // Shared borrow through `ParentRef`: map values are live boxed + // statements (the map owns one intrusive ref on each) and `last_used` + // is a `Cell`. + ParentRef::from(NonNull::new(stmt_ptr).expect("map entries are non-null")) + .last_used + .set(self.statement_lru_clock.get()); + Some(stmt_ptr) + } + + /// Reserve a cache slot for a statement `lookup_statement` just missed on, + /// evicting LRU idle statements first so the cache stays under + /// [`MAX_CACHED_PREPARED_STATEMENTS`]. Returns the raw value-slot pointer; + /// the caller stores the new statement through it. + /// + /// `JsCell::with_mut` scopes the `&mut PreparedStatementsMap` to the + /// `get_or_put` call (single-JS-thread; no re-entry into JS until after + /// the raw value-slot ptr is captured), so callers need no further `&mut` + /// to the map. + pub(crate) fn put_statement( + &self, + name: &[u8], + ) -> Result<*mut *mut PostgresSQLStatement, bun_core::AllocError> { + self.evict_lru_statements(); + self.statements.with_mut(|s| { + s.get_or_put(name) + .map(|e| core::ptr::from_mut::<*mut PostgresSQLStatement>(e.value_ptr)) + }) + } + + /// Evict idle least-recently-used cached statements until the cache is + /// under [`MAX_CACHED_PREPARED_STATEMENTS`], writing a `Close('S', name)` + /// for each so the server drops its side of the named statement. The + /// paired CloseComplete is consumed by `on()` without touching the + /// request queue. + fn evict_lru_statements(&self) { + // `write_bind` may be mid-frame on the stack (a bound value's + // toJSON/valueOf/toString can re-enter `do_run` on this connection): + // appending a Close now would be folded into that Bind's length and + // yield 08P01. Defer to the next insert; a running query's statement + // is not idle anyway, so this only postpones eviction. + if self.has_query_running() { + return; + } + while self.statements.get().len() >= MAX_CACHED_PREPARED_STATEMENTS { + let mut victim: Option> = None; + let mut oldest = u64::MAX; + for value in self.statements.get().values() { + let ptr = NonNull::new(*value).expect("map entries are non-null"); + // Shared borrow; every field read below is `Cell`/`Copy`. + let stmt = ParentRef::from(ptr); + // A `PostgresSQLQuery` still holds this statement (pending, + // running, or not yet finalized): its server-side name must + // stay valid for the Bind that query may still send. + if !stmt.has_one_ref() { + continue; + } + if victim.is_none() || stmt.last_used.get() < oldest { + oldest = stmt.last_used.get(); + victim = Some(ptr); + } + } + // Every cached statement is still referenced by a query; nothing + // can be released right now. The next insert tries again. + let Some(ptr) = victim else { return }; + // `stmt` borrows the statement's own allocation, disjoint from the + // map's, so it stays live across the `with_mut` below. + let stmt = ParentRef::from(ptr); + // Only a statement the server acknowledged (status `Prepared`) + // owns a named server-side object to deallocate; an idle `Failed` + // one never completed a Parse (a rejected Parse is unmapped on + // ErrorResponse), so there is nothing to close for it. + if stmt.status == StatementStatus::Prepared + && (protocol::Close { + p: protocol::PortalOrPreparedStatement::PreparedStatement( + &stmt.signature.prepared_statement_name, + ), + }) + .write(&mut self.writer()) + .is_err() + { + // The Close could not be buffered (OOM): keep the cache + // entry so the server-side statement is not orphaned. + return; + } + // The cache is keyed by `signature.name`, which the statement owns + // a copy of (same invariant the ErrorResponse arm relies on). + let removed = self + .statements + .with_mut(|m| m.remove(&stmt.signature.name[..])); + debug_assert!(removed.is_some(), "victim came from the map"); + // SAFETY: `has_one_ref` above ⇒ the map owned the last ref; + // removing the entry transfers it to us to release here. + unsafe { PostgresSQLStatement::deref(ptr.as_ptr()) }; + } + } } // `Writer.connection` is a @@ -3006,18 +3120,11 @@ impl PostgresSQLConnection { debug!("TODO PortalSuspended"); } MessageType::CloseComplete => { + // The only Close this client sends is the statement cache's + // eviction (`evict_lru_statements`). It is not tied to any + // queued query, so consume the acknowledgement without + // touching the request queue. reader.eat_message(&protocol::CLOSE_COMPLETE)?; - let request = self.current().ok_or(AnyPostgresError::ExpectedRequest)?; - if request.status.get() == QueryStatus::Fail { - return Ok(()); - } - request.on_result( - b"CLOSECOMPLETE", - self.global(), - self.js_value.get().get(), - false, - ); - self.update_ref(); } MessageType::CopyInResponse => { reader.skip_message()?; diff --git a/src/sql_jsc/postgres/PostgresSQLQuery.rs b/src/sql_jsc/postgres/PostgresSQLQuery.rs index 4b0509fa8ea1..5b8a3c98b3d5 100644 --- a/src/sql_jsc/postgres/PostgresSQLQuery.rs +++ b/src/sql_jsc/postgres/PostgresSQLQuery.rs @@ -615,14 +615,11 @@ impl PostgresSQLQuery { .get() .contains(ConnectionFlags::USE_UNNAMED_PREPARED_STATEMENTS) { - // Zero-allocation hit probe: `get_or_put` below boxes the key + // Zero-allocation hit probe: `put_statement` below boxes the key // bytes even when the entry already exists, and a hit (an - // already-prepared named statement) is the steady state. - let existing_stmt = connection - .statements - .get() - .get(&signature.name[..]) - .copied(); + // already-prepared named statement) is the steady state. A hit + // also stamps the entry's LRU clock. + let existing_stmt = connection.lookup_statement(&signature.name); if let Some(stmt_ptr) = existing_stmt { this.statement.set(Some(stmt_ptr)); // Route the `&mut` through the audited `statement_mut()` @@ -688,15 +685,11 @@ impl PostgresSQLQuery { break 'enqueue; } - // `JsCell::with_mut` scopes the `&mut PreparedStatementsMap` to - // the `get_or_put` call (single-JS-thread; no re-entry into JS - // until after the raw value-slot ptr is captured). Extract the - // raw slot ptr while the borrow is live so the remainder of - // this block needs no further `&mut` to the map. - let entry_value_ptr = match connection.statements.with_mut(|s| { - s.get_or_put(&signature.name) - .map(|e| std::ptr::from_mut::<*mut PostgresSQLStatement>(e.value_ptr)) - }) { + // `put_statement` enforces the statement-cache cap (evicting + + // closing LRU idle statements) and hands back the raw slot ptr, + // so the remainder of this block needs no further `&mut` to the + // map. + let entry_value_ptr = match connection.put_statement(&signature.name) { Ok(v) => v, Err(err) => { drop(signature); diff --git a/src/sql_jsc/postgres/PostgresSQLStatement.rs b/src/sql_jsc/postgres/PostgresSQLStatement.rs index a2c4647baba2..2306a9a7629d 100644 --- a/src/sql_jsc/postgres/PostgresSQLStatement.rs +++ b/src/sql_jsc/postgres/PostgresSQLStatement.rs @@ -29,6 +29,10 @@ pub struct PostgresSQLStatement { pub error_response: Option, pub needs_duplicate_check: bool, pub fields_flags: DataCellFlags, + /// LRU stamp from the owning connection's statement clock, bumped each + /// time a later query reuses this cached statement; never-reused + /// statements keep 0 and are evicted first. + pub last_used: Cell, } impl Default for PostgresSQLStatement { @@ -45,6 +49,7 @@ impl Default for PostgresSQLStatement { error_response: None, needs_duplicate_check: true, fields_flags: DataCellFlags::default(), + last_used: Cell::new(0), } } } @@ -83,6 +88,14 @@ impl PostgresSQLStatement { self.ref_count.set(n); } + /// Whether the caller holds the only outstanding ref. The statement cache + /// only evicts statements it is the sole owner of (no `PostgresSQLQuery` + /// left that could still bind to the server-side statement name). + #[inline] + pub(crate) fn has_one_ref(&self) -> bool { + bun_ptr::CellRefCounted::ref_count(self).get() == 1 + } + pub fn check_for_duplicate_fields(&mut self) { if !self.needs_duplicate_check { return; diff --git a/test/js/sql/sql-postgres-statement-cache.test.ts b/test/js/sql/sql-postgres-statement-cache.test.ts new file mode 100644 index 000000000000..56c8911de983 --- /dev/null +++ b/test/js/sql/sql-postgres-statement-cache.test.ts @@ -0,0 +1,351 @@ +// This test counts the exact messages Bun's Postgres client puts on the wire +// (Parse / Close), which needs a scripted server; a real server only exposes +// them through per-session state (pg_prepared_statements) that extra queries +// would perturb. All wire-protocol bytes come from test/js/sql/wire-frames.ts. +// +// Regression test: the client never sent Close and kept an unbounded +// per-connection prepared-statement cache. Every distinct query text on a +// connection kept a named prepared statement allocated in the server session +// and a PostgresSQLStatement (rooting one JSC Structure through a Strong +// handle) on the client, both until the connection closed, so a long-lived +// pooled connection running many distinct query texts (what ORMs produce) +// grew both sides without bound. The client now caps the per-connection cache +// (MAX_CACHED_PREPARED_STATEMENTS in src/sql_jsc/postgres/PostgresSQLConnection.rs) +// and sends Close('S', name) for what it evicts. +import { SQL } from "bun"; +import { heapStats } from "bun:jsc"; +import { expect, test } from "bun:test"; +import { describeWithContainer } from "harness"; +import { + listeningServer, + pgAuthenticationOk, + pgBindComplete, + pgCloseComplete, + pgCommandComplete, + pgDataRow, + pgErrorResponse, + pgParameterDescription, + pgParseComplete, + pgReadCString, + pgReadFrontendMessages, + pgReadyForQuery, + pgRowDescription, +} from "./wire-frames"; + +// Keep in sync with MAX_CACHED_PREPARED_STATEMENTS in +// src/sql_jsc/postgres/PostgresSQLConnection.rs. +const MAX_CACHED_PREPARED_STATEMENTS = 256; + +const OID_TEXT = 25; + +/** + * Minimal Postgres server speaking the extended query protocol: accepts any + * startup, answers Parse / Describe / Bind / Execute / Sync for statements + * that exist, and records every named statement the client prepares (Parse) + * or deallocates (Close 'S'). Like a real server it rejects a Bind to a name + * that was closed (or never prepared) with SQLSTATE 26000 and a Parse that + * redefines a live name with 42P05, then discards messages until Sync; a + * client that closes a statement another query still needs fails loudly. + */ +async function statementCountingServer() { + const counters = { + parses: 0, + executes: 0, + /** named statements the client prepared */ + prepared: new Set(), + /** named statements the client sent Close('S') for */ + closed: new Set(), + }; + const { server, port } = await listeningServer(socket => { + let buffered = Buffer.alloc(0); + let sawStartup = false; + /** named statements currently live in this session */ + const live = new Set(); + /** first parameter of the last Bind, echoed back by the next Execute */ + let lastBoundParam: Buffer | null = null; + // After an ErrorResponse the backend discards messages until Sync. + let skipUntilSync = false; + socket.on("data", chunk => { + buffered = Buffer.concat([buffered, chunk]); + if (!sawStartup) { + // StartupMessage is the one untagged frontend message: Int32(length + // including itself) then the body. + if (buffered.length < 4) return; + const len = buffered.readInt32BE(0); + if (buffered.length < len) return; + buffered = buffered.subarray(len); + sawStartup = true; + socket.write(Buffer.concat([pgAuthenticationOk(), pgReadyForQuery()])); + } + buffered = pgReadFrontendMessages(buffered, (tag, body) => { + if (skipUntilSync && tag !== 0x53 /* Sync */) return; + switch (tag) { + // Parse: String(name) String(query) Int16(nparams) Int32[nparams] + case 0x50: { + const name = pgReadCString(body, 0); + counters.parses++; + if (live.has(name.value)) { + socket.write( + pgErrorResponse({ + S: "ERROR", + C: "42P05", + M: `prepared statement "${name.value}" already exists`, + }), + ); + skipUntilSync = true; + break; + } + live.add(name.value); + counters.prepared.add(name.value); + socket.write(pgParseComplete()); + break; + } + // Describe: Byte1('S' | 'P') String(name) + case 0x44: { + socket.write( + Buffer.concat([pgParameterDescription([OID_TEXT]), pgRowDescription([{ name: "c", typeOid: OID_TEXT }])]), + ); + break; + } + // Bind: String(portal) String(statement) Int16(nformats) Int16[nformats] + // Int16(nparams) (Int32(len) Byte[len])[nparams] ... + case 0x42: { + const portal = pgReadCString(body, 0); + const statement = pgReadCString(body, portal.end); + if (!live.has(statement.value)) { + socket.write( + pgErrorResponse({ + S: "ERROR", + C: "26000", + M: `prepared statement "${statement.value}" does not exist`, + }), + ); + skipUntilSync = true; + break; + } + let offset = statement.end; + const formats = body.readInt16BE(offset); + offset += 2 + 2 * formats; + const params = body.readInt16BE(offset); + offset += 2; + lastBoundParam = null; + if (params > 0) { + const length = body.readInt32BE(offset); + offset += 4; + if (length >= 0) lastBoundParam = Buffer.from(body.subarray(offset, offset + length)); + } + socket.write(pgBindComplete()); + break; + } + // Execute: String(portal) Int32(maxrows) + case 0x45: { + counters.executes++; + socket.write(Buffer.concat([pgDataRow([lastBoundParam]), pgCommandComplete("SELECT 1")])); + break; + } + // Close: Byte1('S' | 'P') String(name). Closing a nonexistent name + // is not an error, but only names the client prepared should ever + // show up here (asserted by the tests through `counters.closed`). + case 0x43: { + if (body[0] === 0x53 /* 'S' */) { + const name = pgReadCString(body, 1); + live.delete(name.value); + counters.closed.add(name.value); + } + socket.write(pgCloseComplete()); + break; + } + // Sync + case 0x53: { + skipUntilSync = false; + socket.write(pgReadyForQuery()); + break; + } + // Terminate + case 0x58: { + socket.end(); + break; + } + } + }); + }); + socket.on("error", () => {}); + }); + return { server, port, counters }; +} + +function sqlUrl(port: number): string { + return `postgres://postgres@127.0.0.1:${port}/postgres`; +} + +// Exceeding the cache cap takes MAX_CACHED_PREPARED_STATEMENTS + 8 distinct +// prepared statements by definition, which outgrows the 5s default per-test +// timeout on debug + ASAN builds, so these tests declare their real budget. +test("postgres: the prepared-statement cache is capped and evicted statements are closed", async () => { + const { server, port, counters } = await statementCountingServer(); + try { + await using sql = new SQL({ url: sqlUrl(port), max: 1 }); + + // More distinct query texts than the cap, all in flight at once on one + // connection: none may be evicted or closed while a query references it + // (the mock rejects a Bind to a closed name, failing the Promise.all). + const DISTINCT = MAX_CACHED_PREPARED_STATEMENTS + 8; + const results = await Promise.all( + Array.from({ length: DISTINCT }, (_, i) => sql.unsafe(`select $1 as c${i}`, [String(i)])), + ); + expect(results).toEqual(Array.from({ length: DISTINCT }, (_, i) => [{ c: String(i) }])); + expect(counters.parses).toBe(DISTINCT); + expect(counters.closed.size).toBe(0); + + // Statements become evictable once their (collected) query wrappers stop + // referencing them. Collect, then keep inserting distinct texts: each one + // past the cap must evict + Close. Poll, since GC sets the pace. + let extra = 0; + while (counters.parses - counters.closed.size > MAX_CACHED_PREPARED_STATEMENTS && extra < 64) { + Bun.gc(true); + await Bun.sleep(0); + await sql.unsafe(`select $1 as extra${extra}`, [String(extra)]); + extra++; + } + + // Without Close the number of live server-side statements + // (parses - closes) grows monotonically with every distinct query text; + // the loop above then exhausts its budget with closed.size still 0. + expect(counters.closed.size).toBeGreaterThan(0); + expect(counters.parses - counters.closed.size).toBeLessThanOrEqual(MAX_CACHED_PREPARED_STATEMENTS); + // Only names the server actually saw prepared may be closed. + expect([...counters.closed].every(name => counters.prepared.has(name))).toBe(true); + // Every query text was parsed and executed exactly once (a Bind to a + // closed name would have been rejected by the mock and thrown above). + expect(counters.parses).toBe(DISTINCT + extra); + expect(counters.executes).toBe(DISTINCT + extra); + } finally { + await new Promise(resolve => server.close(() => resolve())); + } +}, 30_000); + +// The same unbounded cache also leaked client memory: every cached statement +// roots one JSC Structure (its row shape) through a Strong handle, plus its +// own metadata, for the connection's lifetime. Eviction must release both. +test("postgres: evicting cached statements releases their rooted row Structures (client memory)", async () => { + const protectedStructures = () => heapStats().protectedObjectTypeCounts.Structure ?? 0; + const { server, port, counters } = await statementCountingServer(); + try { + await using sql = new SQL({ url: sqlUrl(port), max: 1 }); + + // The mock names every result column "c" and echoes the first bound + // parameter, which also proves the DataRow framing is consumed. + expect(await sql.unsafe("select $1 as warmup", ["w"])).toEqual([{ c: "w" }]); + Bun.gc(true); + await Bun.sleep(0); + const baseline = protectedStructures(); + + const DISTINCT = MAX_CACHED_PREPARED_STATEMENTS + 8; + for (let i = 0; i < DISTINCT; i++) { + await sql.unsafe(`select $1 as c${i}`, [String(i)]); + } + + // Statements become evictable once their collected query wrappers drop + // them (same convergence loop as the capped-eviction test above). + let extra = 0; + while (counters.parses - counters.closed.size > MAX_CACHED_PREPARED_STATEMENTS && extra < 64) { + Bun.gc(true); + await Bun.sleep(0); + await sql.unsafe(`select $1 as extra${extra}`, [String(extra)]); + extra++; + } + Bun.gc(true); + await Bun.sleep(0); + + expect(counters.closed.size).toBeGreaterThan(0); + // Post-GC, only statements still cached may root a Structure (the slack + // absorbs unrelated Strong handles created while the test runs). Without + // eviction every one of the DISTINCT + extra texts roots its Structure + // (and retains its statement) until the connection closes. + expect(protectedStructures() - baseline).toBeLessThanOrEqual(MAX_CACHED_PREPARED_STATEMENTS + 16); + } finally { + await new Promise(resolve => server.close(() => resolve())); + } +}, 30_000); + +// The scripted server above observes the wire; a real server additionally +// proves the Close is accepted where the client writes it (between two +// extended-query sequences) and that pg_prepared_statements, the session's +// source of truth, stays bounded. +describeWithContainer("postgres: statement cache against a real server", { image: "postgres_plain" }, container => { + test("pg_prepared_statements stays within the cache cap", async () => { + await container.ready; + await using sql = new SQL({ + db: "bun_sql_test", + username: "bun_sql_test", + host: container.host, + port: container.port, + max: 1, + }); + // Simple-protocol query: observes the session's named statements without + // creating one. + const statementCount = async () => + Number((await sql`select count(*)::int as n from pg_prepared_statements`.simple())[0].n); + // A statement is only evictable once its (collected) query wrapper stops + // referencing it, and eviction only runs when a new distinct text is + // inserted, so converge by collecting + inserting. `tag` must not repeat + // across calls: only a text that is not already cached triggers eviction. + const convergeUnderCap = async (tag: string) => { + let extra = 0; + while ((await statementCount()) > MAX_CACHED_PREPARED_STATEMENTS && extra < 64) { + Bun.gc(true); + await Bun.sleep(0); + await sql.unsafe(`select $1::text as ${tag}${extra}`, [String(extra)]); + extra++; + } + return statementCount(); + }; + + const DISTINCT = MAX_CACHED_PREPARED_STATEMENTS + 44; + for (let i = 0; i < DISTINCT; i++) { + expect(await sql.unsafe(`select $1::text as c${i}`, [String(i)])).toEqual([{ [`c${i}`]: String(i) }]); + } + + const settled = await convergeUnderCap("extra"); + expect(settled).toBeGreaterThan(0); + expect(settled).toBeLessThanOrEqual(MAX_CACHED_PREPARED_STATEMENTS); + + // Texts whose statements were evicted re-prepare transparently. The count + // exceeds the cap again while their (not yet collected) query wrappers pin + // them, so it must re-converge, proving re-prepared entries are evictable. + for (let i = 0; i < DISTINCT; i++) { + expect(await sql.unsafe(`select $1::text as c${i}`, [String(i)])).toEqual([{ [`c${i}`]: String(i) }]); + } + expect(await convergeUnderCap("again")).toBeLessThanOrEqual(MAX_CACHED_PREPARED_STATEMENTS); + }, 40_000); +}); + +test("postgres: an identical query text keeps reusing one prepared statement and is never closed", async () => { + const { server, port, counters } = await statementCountingServer(); + try { + await using sql = new SQL({ url: sqlUrl(port), max: 1 }); + + // Same text + same parameter type = same statement-cache entry. Collect in + // between so finalized query wrappers cannot take the cached statement (or + // its server-side name) with them. + for (let i = 0; i < 50; i++) { + expect(await sql.unsafe("select $1 as reused", [String(i)])).toEqual([{ c: String(i) }]); + if (i % 10 === 0) { + Bun.gc(true); + await Bun.sleep(0); + } + } + + expect({ + parses: counters.parses, + executes: counters.executes, + closed: counters.closed.size, + }).toEqual({ + parses: 1, + executes: 50, + closed: 0, + }); + } finally { + await new Promise(resolve => server.close(() => resolve())); + } +}); diff --git a/test/js/sql/wire-frames.ts b/test/js/sql/wire-frames.ts index bc452192850a..b0ec34d96f11 100644 --- a/test/js/sql/wire-frames.ts +++ b/test/js/sql/wire-frames.ts @@ -133,6 +133,17 @@ export function pgCommandComplete(tag: string): Buffer { return pgRaw("C", Buffer.concat([Buffer.from(tag), Buffer.from([0])])); } +// PostgreSQL FE/BE protocol §55.7 CloseComplete: Byte1('3') Int32(4) +export function pgCloseComplete(): Buffer { + return pgRaw("3", Buffer.alloc(0)); +} + +/** Read the NUL-terminated String at `offset` in a frontend message body. `end` is the index past the NUL. */ +export function pgReadCString(body: Buffer, offset: number): { value: string; end: number } { + const nul = body.indexOf(0, offset); + return { value: body.toString("utf-8", offset, nul), end: nul + 1 }; +} + export type PgRowDescriptionColumn = { name: string; tableOid?: number;