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/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/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 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)`