From bc8838a461847fdc7c9e46d31ce4a38846350bbe Mon Sep 17 00:00:00 2001 From: Mirhet Julardzija Date: Wed, 2 Sep 2026 23:16:50 +0200 Subject: [PATCH 1/2] Send transaction control on the simple protocol With `prepare: true`, `begin()` sent `savepoint`, `rollback to`, `rollback`, `commit` and `prepare transaction` as tagged templates, so each became a named prepared statement cached on the client connection. On a transaction pooler that does not track named statements (or after the pooler evicts one), the Bind of a named `commit` reaches a backend that never parsed it and fails with SQLSTATE 26000. That error aborts the transaction, and `FetchPreparedStatement` is in `retryRoutines`, so the driver re-sends `commit` on the aborted transaction. Postgres answers that COMMIT with ROLLBACK and no error, and the transaction's writes are lost silently. Reproduced against a Supavisor pooler in transaction mode with one client connection and a second client keeping the other backends busy, twenty sequential `sql.begin` inserts, body statement unnamed: before: 9 of 20 rows, 31 `commit` sends, 0 errors after: 20 of 20 rows, 20 `commit` sends, 0 errors `begin` already goes through `unsafe`. This sends the other five control statements the same way. A zero-argument `unsafe` is `simple: true`, so no `prepare` setting can name it, and the transaction lifecycle (connection close, release only when idle, per-query error capture) is unchanged. Identifier quoting matches `sql(name)`; the `prepare transaction` name doubles single quotes instead of being spliced in raw. --- cf/src/index.js | 18 +++++++++++++----- cjs/src/index.js | 18 +++++++++++++----- deno/src/index.js | 18 +++++++++++++----- src/index.js | 18 +++++++++++++----- 4 files changed, 52 insertions(+), 20 deletions(-) diff --git a/cf/src/index.js b/cf/src/index.js index ffbe7ae..e7dbf80 100644 --- a/cf/src/index.js +++ b/cf/src/index.js @@ -232,6 +232,14 @@ function Postgres(a, b) { } } + // Transaction control runs on the simple protocol (a zero-argument unsafe) so that + // `prepare: true` never names it. A named `commit` whose Bind reaches a pooler backend + // that never parsed it fails with 26000, and the retry of that error re-sends `commit` + // on the now aborted transaction, which Postgres answers with ROLLBACK and no error. + function quoteIdent(name) { + return '"' + name.replace(/"/g, '""') + '"' + } + async function begin(options, fn) { !fn && (fn = options, options = '') const queries = Queue() @@ -256,7 +264,7 @@ function Postgres(a, b) { let uncaughtError , result - name && await sql`savepoint ${ sql(name) }` + name && await sql.unsafe('savepoint ' + quoteIdent(name)) try { result = await new Promise((resolve, reject) => { const x = fn(sql) @@ -267,16 +275,16 @@ function Postgres(a, b) { throw uncaughtError } catch (e) { await (name - ? sql`rollback to ${ sql(name) }` - : sql`rollback` + ? sql.unsafe('rollback to ' + quoteIdent(name)) + : sql.unsafe('rollback') ) throw e instanceof PostgresError && e.code === '25P02' && uncaughtError || e } if (!name) { prepare - ? await sql`prepare transaction '${ sql.unsafe(prepare) }'` - : await sql`commit` + ? await sql.unsafe('prepare transaction \'' + prepare.replace(/'/g, '\'\'') + '\'') + : await sql.unsafe('commit') } return result diff --git a/cjs/src/index.js b/cjs/src/index.js index f09c61c..f5e7076 100644 --- a/cjs/src/index.js +++ b/cjs/src/index.js @@ -231,6 +231,14 @@ function Postgres(a, b) { } } + // Transaction control runs on the simple protocol (a zero-argument unsafe) so that + // `prepare: true` never names it. A named `commit` whose Bind reaches a pooler backend + // that never parsed it fails with 26000, and the retry of that error re-sends `commit` + // on the now aborted transaction, which Postgres answers with ROLLBACK and no error. + function quoteIdent(name) { + return '"' + name.replace(/"/g, '""') + '"' + } + async function begin(options, fn) { !fn && (fn = options, options = '') const queries = Queue() @@ -255,7 +263,7 @@ function Postgres(a, b) { let uncaughtError , result - name && await sql`savepoint ${ sql(name) }` + name && await sql.unsafe('savepoint ' + quoteIdent(name)) try { result = await new Promise((resolve, reject) => { const x = fn(sql) @@ -266,16 +274,16 @@ function Postgres(a, b) { throw uncaughtError } catch (e) { await (name - ? sql`rollback to ${ sql(name) }` - : sql`rollback` + ? sql.unsafe('rollback to ' + quoteIdent(name)) + : sql.unsafe('rollback') ) throw e instanceof PostgresError && e.code === '25P02' && uncaughtError || e } if (!name) { prepare - ? await sql`prepare transaction '${ sql.unsafe(prepare) }'` - : await sql`commit` + ? await sql.unsafe('prepare transaction \'' + prepare.replace(/'/g, '\'\'') + '\'') + : await sql.unsafe('commit') } return result diff --git a/deno/src/index.js b/deno/src/index.js index b6d23db..3977d87 100644 --- a/deno/src/index.js +++ b/deno/src/index.js @@ -232,6 +232,14 @@ function Postgres(a, b) { } } + // Transaction control runs on the simple protocol (a zero-argument unsafe) so that + // `prepare: true` never names it. A named `commit` whose Bind reaches a pooler backend + // that never parsed it fails with 26000, and the retry of that error re-sends `commit` + // on the now aborted transaction, which Postgres answers with ROLLBACK and no error. + function quoteIdent(name) { + return '"' + name.replace(/"/g, '""') + '"' + } + async function begin(options, fn) { !fn && (fn = options, options = '') const queries = Queue() @@ -256,7 +264,7 @@ function Postgres(a, b) { let uncaughtError , result - name && await sql`savepoint ${ sql(name) }` + name && await sql.unsafe('savepoint ' + quoteIdent(name)) try { result = await new Promise((resolve, reject) => { const x = fn(sql) @@ -267,16 +275,16 @@ function Postgres(a, b) { throw uncaughtError } catch (e) { await (name - ? sql`rollback to ${ sql(name) }` - : sql`rollback` + ? sql.unsafe('rollback to ' + quoteIdent(name)) + : sql.unsafe('rollback') ) throw e instanceof PostgresError && e.code === '25P02' && uncaughtError || e } if (!name) { prepare - ? await sql`prepare transaction '${ sql.unsafe(prepare) }'` - : await sql`commit` + ? await sql.unsafe('prepare transaction \'' + prepare.replace(/'/g, '\'\'') + '\'') + : await sql.unsafe('commit') } return result diff --git a/src/index.js b/src/index.js index c7fba3d..5882dc3 100644 --- a/src/index.js +++ b/src/index.js @@ -231,6 +231,14 @@ function Postgres(a, b) { } } + // Transaction control runs on the simple protocol (a zero-argument unsafe) so that + // `prepare: true` never names it. A named `commit` whose Bind reaches a pooler backend + // that never parsed it fails with 26000, and the retry of that error re-sends `commit` + // on the now aborted transaction, which Postgres answers with ROLLBACK and no error. + function quoteIdent(name) { + return '"' + name.replace(/"/g, '""') + '"' + } + async function begin(options, fn) { !fn && (fn = options, options = '') const queries = Queue() @@ -255,7 +263,7 @@ function Postgres(a, b) { let uncaughtError , result - name && await sql`savepoint ${ sql(name) }` + name && await sql.unsafe('savepoint ' + quoteIdent(name)) try { result = await new Promise((resolve, reject) => { const x = fn(sql) @@ -266,16 +274,16 @@ function Postgres(a, b) { throw uncaughtError } catch (e) { await (name - ? sql`rollback to ${ sql(name) }` - : sql`rollback` + ? sql.unsafe('rollback to ' + quoteIdent(name)) + : sql.unsafe('rollback') ) throw e instanceof PostgresError && e.code === '25P02' && uncaughtError || e } if (!name) { prepare - ? await sql`prepare transaction '${ sql.unsafe(prepare) }'` - : await sql`commit` + ? await sql.unsafe('prepare transaction \'' + prepare.replace(/'/g, '\'\'') + '\'') + : await sql.unsafe('commit') } return result From 90062b17c79d67d35d308739bb8a0114cd12c47b Mon Sep 17 00:00:00 2001 From: Chris Bala Date: Tue, 29 Sep 2026 19:43:41 -0700 Subject: [PATCH 2/2] Do not retry a prepared statement that failed inside a transaction A prepared statement that fails because its cached plan went stale (RevalidateCachedQuery) or it no longer exists (FetchPreparedStatement) is prepared again and retried. Inside a transaction block that retry cannot work: the error has already aborted the transaction, and a query pipelined behind the failed one may already have ended it. - With concurrent queries in `sql.begin()`, the retry gets 25P02, the answers fall out of step, ROLLBACK is never completed, and the connection is never released (#1234). - With the transaction pipelined on a reserved connection, the retry runs after `commit` has ended the transaction (as ROLLBACK), on its own, so the retried write lands while the rest of the transaction does not: const reserved = await sql.reserve() await Promise.all([ reserved`begin`, reserved`insert into test values (2, 2)`, reserved`update test set a = a + ${ 1 } where id = 1 returning *`, reserved`commit` ]) // after a concurrent `alter table test add column b int`: the insert // is rolled back, but the update lands after it, outside the // transaction Retry only when ReadyForQuery reports the connection idle, that is, not in a transaction block. Otherwise reject the query with its original error, and forget its statement so that its next use prepares it again. This depends on #1212, which this branch is based on. Before it, `sql.begin()` sent its own `rollback` and `commit` as prepared statements, and a `rollback` whose statement had been deallocated only succeeded through this retry (the test "Properly throws routine error on not prepared statements in transaction" fails without #1212). Fixes #1234 Co-Authored-By: Claude Opus 5.5 (1M context) --- src/connection.js | 12 +++++++++++- tests/index.js | 37 +++++++++++++++++++++++++++++++++++++ 2 files changed, 48 insertions(+), 1 deletion(-) diff --git a/src/connection.js b/src/connection.js index 10ab1bb..40018ae 100644 --- a/src/connection.js +++ b/src/connection.js @@ -538,7 +538,9 @@ function Connection(options, queues = {}, { onopen = noop, onend = noop, onclose query.retried ? errored(query.retried) : query.prepared && retryRoutines.has(errorResponse.routine) - ? retry(query, errorResponse) + ? x[5] === 73 // I: not in a transaction block + ? retry(query, errorResponse) + : forget(query, errorResponse) : errored(errorResponse) } else { query.resolve(results || result) @@ -824,6 +826,14 @@ function Connection(options, queues = {}, { onopen = noop, onend = noop, onclose execute(q) } + // A failed transaction ignores every command until it ends, and a query + // pipelined behind this one may already be ending it, so the query is not + // retried. Its stale statement is forgotten, to be prepared anew next time. + function forget(q, error) { + delete statements[q.signature] + errored(error) + } + function NotificationResponse(x) { if (!onnotify) return diff --git a/tests/index.js b/tests/index.js index 845c6ce..1094b4f 100644 --- a/tests/index.js +++ b/tests/index.js @@ -1892,6 +1892,43 @@ t('Does not re-serialize bytea parameters when retrying on RevalidateCachedQuery ] }) +t('Stale prepared statements in a transaction returning an array of queries reject and release the connection', async() => { + const other = postgres(options) + await sql`create table test (id int, a int)` + await sql`insert into test values (1, 1)` + await other`select * from test where id = ${ 1 }` + await other`select id from test where id = ${ 1 } and a > ${ 0 }` + await sql`alter table test add column b int` + const error = await other.begin(sql => [ + sql`select * from test where id = ${ 1 }`, + sql`select id from test where id = ${ 1 } and a > ${ 0 }` + ]).catch(e => e.code) + const [row] = await other`select * from test where id = ${ 1 }` + await sql`drop table test` + await other.end() + return ['0A000,true', [error, 'b' in row].join()] +}) + +t('Stale prepared statement is not retried after the transaction it was pipelined in ends', async() => { + const other = postgres(options) + await sql`create table test (id int, a int)` + await sql`insert into test values (1, 1)` + await other`update test set a = a + ${ 1 } where id = 1 returning *` + await sql`alter table test add column b int` + const reserved = await other.reserve() + const results = await Promise.all([ + reserved`begin`, + reserved`insert into test values (2, 2)`, + reserved`update test set a = a + ${ 1 } where id = 1 returning *`, + reserved`commit` + ].map(x => x.then(x => x.command, e => e.code))) + reserved.release() + const rows = (await sql`select id, a from test order by id`).map(x => x.id + ':' + x.a) + await sql`drop table test` + await other.end() + return ['BEGIN,INSERT,0A000,ROLLBACK,1:2', [...results, ...rows].join()] +}) + t('Does not re-serialize bytea parameters when retrying on FetchPreparedStatement error', async() => { const insert = () => sql`insert into test (data) values (${ Buffer.from('hello') })` await sql`create table test (data bytea)`