diff --git a/cf/src/connection.js b/cf/src/connection.js index 8e79170a..833094d6 100644 --- a/cf/src/connection.js +++ b/cf/src/connection.js @@ -562,8 +562,9 @@ function Connection(options, queues = {}, { onopen = noop, onend = noop, onclose } if (needsTypes) { - initial.reserve && (initial = null) - return fetchArrayTypes() + const reserve = initial.reserve ? initial : null + reserve && (initial = null) + return fetchArrayTypes(reserve) } initial && !initial.reserve && execute(initial) @@ -572,6 +573,9 @@ function Connection(options, queues = {}, { onopen = noop, onend = noop, onclose return } + if (needsTypes) + return terminate() + while (sent.length && (query = sent.shift()) && (query.active = true, query.cancelled)) Connection(options).cancel(query.state, query.cancelled.resolve, query.cancelled.reject) @@ -767,9 +771,9 @@ function Connection(options, queues = {}, { onopen = noop, onend = noop, onclose backend.secret = x.readUInt32BE(9) } - async function fetchArrayTypes() { + async function fetchArrayTypes(reserve) { needsTypes = false - const types = await new Query([` + const query = new Query([` select b.oid, b.typarray from pg_catalog.pg_type a left join pg_catalog.pg_type b on b.oid = a.typelem @@ -777,7 +781,11 @@ function Connection(options, queues = {}, { onopen = noop, onend = noop, onclose group by b.oid, b.typarray order by b.oid `], [], execute) - types.forEach(({ oid, typarray }) => addArrayType(oid, typarray)) + const resolve = query.resolve + query.resolve = types => (types.forEach(({ oid, typarray }) => addArrayType(oid, typarray)), resolve(types)) + const reject = query.reject + query.reject = err => (needsTypes = true, reserve && reserve.reject(err), reject(err)) + await query.catch(() => { /* settled by query.reject above; ReadyForQuery closes the connection */ }) } function addArrayType(oid, typarray) { diff --git a/cjs/src/connection.js b/cjs/src/connection.js index 07f67167..a29bf703 100644 --- a/cjs/src/connection.js +++ b/cjs/src/connection.js @@ -560,8 +560,9 @@ function Connection(options, queues = {}, { onopen = noop, onend = noop, onclose } if (needsTypes) { - initial.reserve && (initial = null) - return fetchArrayTypes() + const reserve = initial.reserve ? initial : null + reserve && (initial = null) + return fetchArrayTypes(reserve) } initial && !initial.reserve && execute(initial) @@ -570,6 +571,9 @@ function Connection(options, queues = {}, { onopen = noop, onend = noop, onclose return } + if (needsTypes) + return terminate() + while (sent.length && (query = sent.shift()) && (query.active = true, query.cancelled)) Connection(options).cancel(query.state, query.cancelled.resolve, query.cancelled.reject) @@ -765,9 +769,9 @@ function Connection(options, queues = {}, { onopen = noop, onend = noop, onclose backend.secret = x.readUInt32BE(9) } - async function fetchArrayTypes() { + async function fetchArrayTypes(reserve) { needsTypes = false - const types = await new Query([` + const query = new Query([` select b.oid, b.typarray from pg_catalog.pg_type a left join pg_catalog.pg_type b on b.oid = a.typelem @@ -775,7 +779,11 @@ function Connection(options, queues = {}, { onopen = noop, onend = noop, onclose group by b.oid, b.typarray order by b.oid `], [], execute) - types.forEach(({ oid, typarray }) => addArrayType(oid, typarray)) + const resolve = query.resolve + query.resolve = types => (types.forEach(({ oid, typarray }) => addArrayType(oid, typarray)), resolve(types)) + const reject = query.reject + query.reject = err => (needsTypes = true, reserve && reserve.reject(err), reject(err)) + await query.catch(() => { /* settled by query.reject above; ReadyForQuery closes the connection */ }) } function addArrayType(oid, typarray) { diff --git a/cjs/tests/index.js b/cjs/tests/index.js index 85d1aa46..5e4f431c 100644 --- a/cjs/tests/index.js +++ b/cjs/tests/index.js @@ -154,6 +154,69 @@ t('Escape in arrays', async() => ['Hello "you",c:\\windows', (await sql`select ${ sql.array(['Hello "you"', 'c:\\windows']) } as x`)[0].x.join(',')] ) +t('Array in first query of new client', async() => { + const sql = postgres(options) + const x = (await sql`select ${ sql.array([1, 2]) }::int[] as x`)[0].x + await sql.end() + return ['1,2', x.join()] +}) + +async function withoutArrayTypes(fn) { + const unhandled = [] + , onUnhandled = e => unhandled.push(e.code) + , grant = () => exec('psql', ['-d', 'postgres_js_test', '-c', 'grant select on pg_catalog.pg_type to public']) + + // Deno's process.on('unhandledRejection') is a stub, but Deno exits on an unhandled rejection anyway + globalThis.Deno || process.on('unhandledRejection', onUnhandled) + await exec('psql', ['-c', 'drop user if exists postgres_js_test_no_types']) + await exec('psql', ['-c', 'create user postgres_js_test_no_types']) + await exec('psql', ['-d', 'postgres_js_test', '-c', 'revoke select on pg_catalog.pg_type from public']) + + const sql = postgres({ + ...options, + user: 'postgres_js_test_no_types', + host: process.env.PGSOCKET || '/tmp' // eslint-disable-line + }) + + try { + const result = await fn({ sql, grant }) + await delay(0) // Node reports a rejection left unhandled before the next timer + return JSON.stringify({ ...result, unhandled }) + } finally { + await grant() + globalThis.Deno || process.off('unhandledRejection', onUnhandled) + await sql.end({ timeout: 0 }) + await exec('psql', ['-c', 'drop user postgres_js_test_no_types']) + } +} + +t('Failed array types fetch rejects the first query and keeps the client usable', async() => [ + '{"first":"42501","later":"1,2","unhandled":[]}', + await withoutArrayTypes(async({ sql, grant }) => { + const first = await sql`select 1`.catch(e => e.code) + await grant() + const later = await sql`select ${ sql.array([1, 2]) }::int[] as x`.then(([{ x }]) => x.join(), e => e.code) + return { first, later } + }) +]) + +t('Failed array types fetch rejects every query waiting on the connection', async() => [ + '{"waiting":["42501","42501"],"unhandled":[]}', + await withoutArrayTypes(async({ sql }) => ({ + waiting: await Promise.all([ + sql`select 1`.catch(e => e.code), + sql`select ${ sql.array([3, 4]) }::int[] as x`.catch(e => e.code) + ]) + })) +]) + +t('Failed array types fetch rejects a waiting reserve', async() => [ + '{"reserved":"42501","unhandled":[]}', + await withoutArrayTypes(async({ sql }) => ({ + reserved: await sql.reserve().then(r => (r.release(), 'reserved'), e => e.code) + })) +]) + t('Escapes', async() => { return ['hej"hej', Object.keys((await sql`select 1 as ${ sql('hej"hej') }`)[0])[0]] }) diff --git a/deno/src/connection.js b/deno/src/connection.js index 796725de..0e4367ec 100644 --- a/deno/src/connection.js +++ b/deno/src/connection.js @@ -563,8 +563,9 @@ function Connection(options, queues = {}, { onopen = noop, onend = noop, onclose } if (needsTypes) { - initial.reserve && (initial = null) - return fetchArrayTypes() + const reserve = initial.reserve ? initial : null + reserve && (initial = null) + return fetchArrayTypes(reserve) } initial && !initial.reserve && execute(initial) @@ -573,6 +574,9 @@ function Connection(options, queues = {}, { onopen = noop, onend = noop, onclose return } + if (needsTypes) + return terminate() + while (sent.length && (query = sent.shift()) && (query.active = true, query.cancelled)) Connection(options).cancel(query.state, query.cancelled.resolve, query.cancelled.reject) @@ -768,9 +772,9 @@ function Connection(options, queues = {}, { onopen = noop, onend = noop, onclose backend.secret = x.readUInt32BE(9) } - async function fetchArrayTypes() { + async function fetchArrayTypes(reserve) { needsTypes = false - const types = await new Query([` + const query = new Query([` select b.oid, b.typarray from pg_catalog.pg_type a left join pg_catalog.pg_type b on b.oid = a.typelem @@ -778,7 +782,11 @@ function Connection(options, queues = {}, { onopen = noop, onend = noop, onclose group by b.oid, b.typarray order by b.oid `], [], execute) - types.forEach(({ oid, typarray }) => addArrayType(oid, typarray)) + const resolve = query.resolve + query.resolve = types => (types.forEach(({ oid, typarray }) => addArrayType(oid, typarray)), resolve(types)) + const reject = query.reject + query.reject = err => (needsTypes = true, reserve && reserve.reject(err), reject(err)) + await query.catch(() => { /* settled by query.reject above; ReadyForQuery closes the connection */ }) } function addArrayType(oid, typarray) { diff --git a/deno/tests/index.js b/deno/tests/index.js index cc2a2518..4a321279 100644 --- a/deno/tests/index.js +++ b/deno/tests/index.js @@ -156,6 +156,69 @@ t('Escape in arrays', async() => ['Hello "you",c:\\windows', (await sql`select ${ sql.array(['Hello "you"', 'c:\\windows']) } as x`)[0].x.join(',')] ) +t('Array in first query of new client', async() => { + const sql = postgres(options) + const x = (await sql`select ${ sql.array([1, 2]) }::int[] as x`)[0].x + await sql.end() + return ['1,2', x.join()] +}) + +async function withoutArrayTypes(fn) { + const unhandled = [] + , onUnhandled = e => unhandled.push(e.code) + , grant = () => exec('psql', ['-d', 'postgres_js_test', '-c', 'grant select on pg_catalog.pg_type to public']) + + // Deno's process.on('unhandledRejection') is a stub, but Deno exits on an unhandled rejection anyway + globalThis.Deno || process.on('unhandledRejection', onUnhandled) + await exec('psql', ['-c', 'drop user if exists postgres_js_test_no_types']) + await exec('psql', ['-c', 'create user postgres_js_test_no_types']) + await exec('psql', ['-d', 'postgres_js_test', '-c', 'revoke select on pg_catalog.pg_type from public']) + + const sql = postgres({ + ...options, + user: 'postgres_js_test_no_types', + host: process.env.PGSOCKET || '/tmp' // eslint-disable-line + }) + + try { + const result = await fn({ sql, grant }) + await delay(0) // Node reports a rejection left unhandled before the next timer + return JSON.stringify({ ...result, unhandled }) + } finally { + await grant() + globalThis.Deno || process.off('unhandledRejection', onUnhandled) + await sql.end({ timeout: 0 }) + await exec('psql', ['-c', 'drop user postgres_js_test_no_types']) + } +} + +t('Failed array types fetch rejects the first query and keeps the client usable', async() => [ + '{"first":"42501","later":"1,2","unhandled":[]}', + await withoutArrayTypes(async({ sql, grant }) => { + const first = await sql`select 1`.catch(e => e.code) + await grant() + const later = await sql`select ${ sql.array([1, 2]) }::int[] as x`.then(([{ x }]) => x.join(), e => e.code) + return { first, later } + }) +]) + +t('Failed array types fetch rejects every query waiting on the connection', async() => [ + '{"waiting":["42501","42501"],"unhandled":[]}', + await withoutArrayTypes(async({ sql }) => ({ + waiting: await Promise.all([ + sql`select 1`.catch(e => e.code), + sql`select ${ sql.array([3, 4]) }::int[] as x`.catch(e => e.code) + ]) + })) +]) + +t('Failed array types fetch rejects a waiting reserve', async() => [ + '{"reserved":"42501","unhandled":[]}', + await withoutArrayTypes(async({ sql }) => ({ + reserved: await sql.reserve().then(r => (r.release(), 'reserved'), e => e.code) + })) +]) + t('Escapes', async() => { return ['hej"hej', Object.keys((await sql`select 1 as ${ sql('hej"hej') }`)[0])[0]] }) diff --git a/src/connection.js b/src/connection.js index 10ab1bb3..d8d0ee30 100644 --- a/src/connection.js +++ b/src/connection.js @@ -560,8 +560,9 @@ function Connection(options, queues = {}, { onopen = noop, onend = noop, onclose } if (needsTypes) { - initial.reserve && (initial = null) - return fetchArrayTypes() + const reserve = initial.reserve ? initial : null + reserve && (initial = null) + return fetchArrayTypes(reserve) } initial && !initial.reserve && execute(initial) @@ -570,6 +571,9 @@ function Connection(options, queues = {}, { onopen = noop, onend = noop, onclose return } + if (needsTypes) + return terminate() + while (sent.length && (query = sent.shift()) && (query.active = true, query.cancelled)) Connection(options).cancel(query.state, query.cancelled.resolve, query.cancelled.reject) @@ -765,9 +769,9 @@ function Connection(options, queues = {}, { onopen = noop, onend = noop, onclose backend.secret = x.readUInt32BE(9) } - async function fetchArrayTypes() { + async function fetchArrayTypes(reserve) { needsTypes = false - const types = await new Query([` + const query = new Query([` select b.oid, b.typarray from pg_catalog.pg_type a left join pg_catalog.pg_type b on b.oid = a.typelem @@ -775,7 +779,11 @@ function Connection(options, queues = {}, { onopen = noop, onend = noop, onclose group by b.oid, b.typarray order by b.oid `], [], execute) - types.forEach(({ oid, typarray }) => addArrayType(oid, typarray)) + const resolve = query.resolve + query.resolve = types => (types.forEach(({ oid, typarray }) => addArrayType(oid, typarray)), resolve(types)) + const reject = query.reject + query.reject = err => (needsTypes = true, reserve && reserve.reject(err), reject(err)) + await query.catch(() => { /* settled by query.reject above; ReadyForQuery closes the connection */ }) } function addArrayType(oid, typarray) { diff --git a/tests/index.js b/tests/index.js index 845c6ce0..1e38f2f2 100644 --- a/tests/index.js +++ b/tests/index.js @@ -154,6 +154,69 @@ t('Escape in arrays', async() => ['Hello "you",c:\\windows', (await sql`select ${ sql.array(['Hello "you"', 'c:\\windows']) } as x`)[0].x.join(',')] ) +t('Array in first query of new client', async() => { + const sql = postgres(options) + const x = (await sql`select ${ sql.array([1, 2]) }::int[] as x`)[0].x + await sql.end() + return ['1,2', x.join()] +}) + +async function withoutArrayTypes(fn) { + const unhandled = [] + , onUnhandled = e => unhandled.push(e.code) + , grant = () => exec('psql', ['-d', 'postgres_js_test', '-c', 'grant select on pg_catalog.pg_type to public']) + + // Deno's process.on('unhandledRejection') is a stub, but Deno exits on an unhandled rejection anyway + globalThis.Deno || process.on('unhandledRejection', onUnhandled) + await exec('psql', ['-c', 'drop user if exists postgres_js_test_no_types']) + await exec('psql', ['-c', 'create user postgres_js_test_no_types']) + await exec('psql', ['-d', 'postgres_js_test', '-c', 'revoke select on pg_catalog.pg_type from public']) + + const sql = postgres({ + ...options, + user: 'postgres_js_test_no_types', + host: process.env.PGSOCKET || '/tmp' // eslint-disable-line + }) + + try { + const result = await fn({ sql, grant }) + await delay(0) // Node reports a rejection left unhandled before the next timer + return JSON.stringify({ ...result, unhandled }) + } finally { + await grant() + globalThis.Deno || process.off('unhandledRejection', onUnhandled) + await sql.end({ timeout: 0 }) + await exec('psql', ['-c', 'drop user postgres_js_test_no_types']) + } +} + +t('Failed array types fetch rejects the first query and keeps the client usable', async() => [ + '{"first":"42501","later":"1,2","unhandled":[]}', + await withoutArrayTypes(async({ sql, grant }) => { + const first = await sql`select 1`.catch(e => e.code) + await grant() + const later = await sql`select ${ sql.array([1, 2]) }::int[] as x`.then(([{ x }]) => x.join(), e => e.code) + return { first, later } + }) +]) + +t('Failed array types fetch rejects every query waiting on the connection', async() => [ + '{"waiting":["42501","42501"],"unhandled":[]}', + await withoutArrayTypes(async({ sql }) => ({ + waiting: await Promise.all([ + sql`select 1`.catch(e => e.code), + sql`select ${ sql.array([3, 4]) }::int[] as x`.catch(e => e.code) + ]) + })) +]) + +t('Failed array types fetch rejects a waiting reserve', async() => [ + '{"reserved":"42501","unhandled":[]}', + await withoutArrayTypes(async({ sql }) => ({ + reserved: await sql.reserve().then(r => (r.release(), 'reserved'), e => e.code) + })) +]) + t('Escapes', async() => { return ['hej"hej', Object.keys((await sql`select 1 as ${ sql('hej"hej') }`)[0])[0]] })