From 2efc4454a0f010a93e59eec75e7c43f09e2a86d4 Mon Sep 17 00:00:00 2001 From: Chris Bala Date: Tue, 29 Sep 2026 19:43:41 -0700 Subject: [PATCH] Answer a query that fails to build in its place in the pipeline When a query throws while it is built (an undefined value, more than 65534 parameters, a serializer that throws), `execute` has already queued it to await an answer, but nothing of it was written. The catch then rejected the connection's current query with the new query's error and left the failed query waiting for an answer that never comes, so every later answer on the connection went to the wrong query: await sql`select 1` const a = sql`select pg_sleep(0.1), 1 as x` , b = sql`select ${ undefined }` , c = sql`select 3 as x` // a rejects with UNDEFINED_VALUE, b with "Cannot set properties of // null (setting 'columns')", c never settles, nor does any later query This is also why `sql.begin(sql => [sql`select 1`, sql`select ${ undefined }`])` never settles (#1082). On a reserved connection with the transaction pipelined (`begin`, insert, failing query, insert, `commit`), both inserts committed. Send a query the server is sure to refuse in the failed query's place. The failed query then gets an answer of its own, every later answer stays with its query, and a transaction it was part of is aborted rather than committed without it. The failed query is rejected with its own error, as a retried query already is. It is no longer marked to be described first or as a cursor, since either would make the server's error send a Sync of its own and shift the next answer. Fixes #1082 Co-Authored-By: Claude Opus 5.5 (1M context) --- src/connection.js | 13 +++++++-- tests/index.js | 71 +++++++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 82 insertions(+), 2 deletions(-) diff --git a/src/connection.js b/src/connection.js index 10ab1bb..f100244 100644 --- a/src/connection.js +++ b/src/connection.js @@ -19,6 +19,7 @@ const Sync = b().S().end() , Flush = b().H().end() , SSLRequest = b().i32(8).i32(80877103).end(8) , ExecuteUnnamed = Buffer.concat([b().E().str(b.N).i32(0).end(), Sync]) + , Unbuildable = b().Q().str('postgres.js: query could not be built' + b.N).end() , DescribeUnnamed = b().D().str('S').str(b.N).end() , noop = () => { /* noop */ } @@ -176,8 +177,16 @@ function Connection(options, queues = {}, { onopen = noop, onend = noop, onclose && sent.length < max_pipeline && (!q.options.onexecute || q.options.onexecute(connection)) } catch (error) { - sent.length === 0 && write(Sync) - errored(error) + // q already has its place in the queue of queries awaiting an answer, + // but nothing of it was written. Send a query the server is sure to + // refuse in its place, so that q gets an answer of its own, every + // later answer stays with its query, and a transaction q was part of + // is aborted instead of committing without it. q is then rejected with + // its own error, as a retried query is (ReadyForQuery). + q.retried = error + q.describeFirst = false + q.cursorFn = null + write(Unbuildable) return true } } diff --git a/tests/index.js b/tests/index.js index 845c6ce..1901395 100644 --- a/tests/index.js +++ b/tests/index.js @@ -342,6 +342,77 @@ t('Undefined values throws', async() => { return ['UNDEFINED_VALUE', error] }) +t('Undefined value in a pipelined query rejects only that query', async() => { + const sql = postgres(options) + await sql`select 1` + const results = await Promise.all([ + sql`select pg_sleep(0.1), 1 as x`, + sql`select ${ undefined } as x`, + sql`select 3 as x` + ].map(x => x.then(([x]) => x.x, e => e.code))) + const after = (await sql`select 4 as x`)[0].x + await sql.end() + return ['1,UNDEFINED_VALUE,3,4', [...results, after].join()] +}) + +t('Undefined value in a transaction returning an array of queries rejects', async() => { + const sql = postgres(options) + await sql`create table test (x int)` + const error = await sql.begin(sql => [ + sql`insert into test values (1)`, + sql`select ${ undefined } as x` + ]).catch(e => e.code) + const count = (await sql`select count(*)::int from test`)[0].count + await sql`drop table test` + await sql.end() + return ['UNDEFINED_VALUE,0', [error, count].join()] +}) + +t('Undefined value in a pipelined cursor rejects only that query', async() => { + const sql = postgres(options) + await sql`select 1` + const results = await Promise.all([ + sql`select pg_sleep(0.1), 1 as x`, + sql`select ${ undefined } as x`.cursor(() => { /* noop */ }), + sql`select 3 as x` + ].map(x => x.then(x => x.length ? x[0].x : x.command, e => e.code))) + const after = (await sql`select 4 as x`)[0].x + await sql.end() + return ['1,UNDEFINED_VALUE,3,4', [...results, after].join()] +}) + +t('Too many parameters in a pipelined query rejects only that query', async() => { + const sql = postgres(options) + await sql`select 1` + const results = await Promise.all([ + sql`select pg_sleep(0.1), 1 as x`, + sql.unsafe('select 2 as x', [...Array(65535).keys()]), + sql`select 3 as x`, + sql`select 4 as x` + ].map(x => x.then(([x]) => x.x, e => e.code))) + const after = (await sql`select 5 as x`)[0].x + await sql.end() + return ['1,MAX_PARAMETERS_EXCEEDED,3,4,5', [...results, after].join()] +}) + +t('Undefined value aborts a pipelined transaction on a reserved connection', async() => { + const sql = postgres(options) + await sql`create table test (x int)` + const reserved = await sql.reserve() + const results = await Promise.all([ + reserved`begin`, + reserved`insert into test values (1)`, + reserved`select ${ undefined } as x`, + reserved`insert into test values (2)`, + reserved`commit` + ].map(x => x.then(x => x.command, e => e.code))) + reserved.release() + const count = (await sql`select count(*)::int from test`)[0].count + await sql`drop table test` + await sql.end() + return ['BEGIN,INSERT,UNDEFINED_VALUE,25P02,ROLLBACK,0', [...results, count].join()] +}) + t('Transform undefined', async() => { const sql = postgres({ ...options, transform: { undefined: null } }) return [null, (await sql`select ${ undefined } as x`)[0].x]