Skip to content
Open
Show file tree
Hide file tree
Changes from all 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
18 changes: 13 additions & 5 deletions cf/src/index.js
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand All @@ -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)
Expand All @@ -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
Expand Down
18 changes: 13 additions & 5 deletions cjs/src/index.js
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand All @@ -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)
Expand All @@ -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
Expand Down
18 changes: 13 additions & 5 deletions deno/src/index.js
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand All @@ -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)
Expand All @@ -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
Expand Down
12 changes: 11 additions & 1 deletion src/connection.js
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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
Expand Down
18 changes: 13 additions & 5 deletions src/index.js
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand All @@ -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)
Expand All @@ -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
Expand Down
37 changes: 37 additions & 0 deletions tests/index.js
Original file line number Diff line number Diff line change
Expand Up @@ -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)`
Expand Down