|
| 1 | +'use strict' |
| 2 | + |
| 3 | +// pg against itself: the same statements, on one connection per mode, and every mode must |
| 4 | +// answer what a plain query on a plain client answered. |
| 5 | +// |
| 6 | +// The reference arm is `client.query(text, values)`. The others are the paths a program can |
| 7 | +// take to the same rows: the extended protocol forced on a query that would have used the |
| 8 | +// simple one, a named prepared statement run twice, rowMode array, the binary result format, |
| 9 | +// pipeline mode with the whole round in flight at once, pg-cursor reading in batches and |
| 10 | +// pg-query-stream. Each keeps its own connection and runs the whole round in order, so the |
| 11 | +// temp table and the transaction state are the same on every arm. |
| 12 | +// |
| 13 | +// node fuzz/modes.js fifty rounds, PG* variables say where the server is |
| 14 | +// node fuzz/modes.js --seed 12345 replay what a past run did |
| 15 | +// node fuzz/modes.js --keep-going do not stop at the first divergence |
| 16 | + |
| 17 | +const { createHash } = require('crypto') |
| 18 | +const pg = require('../packages/pg') |
| 19 | +const Cursor = require('../packages/pg-cursor') |
| 20 | +const QueryStream = require('../packages/pg-query-stream') |
| 21 | +const { int, mulberry32, main } = require('./lib') |
| 22 | +const { draw, render, summarize, QUIET, rowsOnly, source, variants, SETUP } = require('./queries') |
| 23 | + |
| 24 | +// the outcome of a query, whichever way it went |
| 25 | +const outcome = (promise) => |
| 26 | + promise.then( |
| 27 | + (result) => ({ result }), |
| 28 | + (error) => ({ error }) |
| 29 | + ) |
| 30 | +const settle = (promise) => outcome(promise).then(({ result, error }) => summarize(result, error)) |
| 31 | +const settleRows = (promise) => outcome(promise).then(({ result, error }) => rowsOnly(result, error)) |
| 32 | + |
| 33 | +const name = (text) => `fuzz_${createHash('sha1').update(text).digest('hex')}` |
| 34 | + |
| 35 | +// each arm runs one query of the round on its own connection and answers the canonical text. |
| 36 | +// `rng` is the round's, so a batch size is drawn the same on a replay |
| 37 | +const ARMS = { |
| 38 | + extended: (client, q) => settle(client.query({ text: q.text, values: q.values, queryMode: 'extended' })), |
| 39 | + // twice, so the second run goes through the statement cache; only a select, since running a |
| 40 | + // write twice would leave this arm's table different from the others |
| 41 | + prepared: async (client, q) => { |
| 42 | + const first = await settle(client.query({ text: q.text, values: q.values, name: name(q.text) })) |
| 43 | + // and not after a failure either, which would have aborted an open transaction |
| 44 | + if (q.kind !== 'select' || first.startsWith('{"error"')) return first |
| 45 | + const second = await settle(client.query({ text: q.text, values: q.values, name: name(q.text) })) |
| 46 | + return first === second ? first : `first run: ${first}\n second run: ${second}` |
| 47 | + }, |
| 48 | + array: async (client, q, rng, reference) => { |
| 49 | + const got = await settle(client.query({ text: q.text, values: q.values, rowMode: 'array' })) |
| 50 | + // the reference rows as arrays, which is the only thing this mode changes |
| 51 | + const { result, error } = reference |
| 52 | + const expected = summarize(result && { ...result, rows: result.rows.map(Object.values) }, error) |
| 53 | + return got === expected ? summarize(result, error) : got |
| 54 | + }, |
| 55 | + binary: (client, q) => |
| 56 | + q.binary |
| 57 | + ? settle(client.query({ text: q.text, values: q.values, binary: true })) |
| 58 | + : settle(client.query(q.text, q.values)), |
| 59 | + cursor: async (client, q, rng) => { |
| 60 | + if (q.kind !== 'select') return settleRows(client.query(q.text, q.values)) |
| 61 | + const cursor = client.query(new Cursor(q.text, q.values)) |
| 62 | + const rows = [] |
| 63 | + try { |
| 64 | + for (;;) { |
| 65 | + const batch = await cursor.read(int(rng, 1, 40)) |
| 66 | + if (batch.length === 0) break |
| 67 | + rows.push(...batch) |
| 68 | + } |
| 69 | + await cursor.close() |
| 70 | + return rowsOnly({ rows }) |
| 71 | + } catch (error) { |
| 72 | + return rowsOnly(null, error) |
| 73 | + } |
| 74 | + }, |
| 75 | + stream: async (client, q, rng) => { |
| 76 | + if (q.kind !== 'select') return settleRows(client.query(q.text, q.values)) |
| 77 | + const rows = [] |
| 78 | + try { |
| 79 | + const stream = client.query( |
| 80 | + new QueryStream(q.text, q.values, { batchSize: int(rng, 1, 40), highWaterMark: int(rng, 1, 40) }) |
| 81 | + ) |
| 82 | + for await (const row of stream) rows.push(row) |
| 83 | + return rowsOnly({ rows }) |
| 84 | + } catch (error) { |
| 85 | + return rowsOnly(null, error) |
| 86 | + } |
| 87 | + }, |
| 88 | +} |
| 89 | + |
| 90 | +const clients = {} |
| 91 | +const connect = async () => { |
| 92 | + clients.reference = new pg.Client() |
| 93 | + clients.pipeline = new pg.Client({ pipeline: true }) |
| 94 | + for (const arm of Object.keys(ARMS)) clients[arm] = new pg.Client() |
| 95 | + for (const client of Object.values(clients)) { |
| 96 | + await client.connect() |
| 97 | + await client.query(QUIET) |
| 98 | + } |
| 99 | +} |
| 100 | + |
| 101 | +// every round starts from the same state on every arm: no transaction open, an empty table |
| 102 | +const reset = async (client) => { |
| 103 | + await client.query('ROLLBACK').catch(() => {}) |
| 104 | + await client.query('DROP TABLE IF EXISTS fuzz_rows') |
| 105 | + await client.query(SETUP) |
| 106 | +} |
| 107 | + |
| 108 | +const run = async (plan) => { |
| 109 | + if (!clients.reference) await connect() |
| 110 | + for (const client of Object.values(clients)) await reset(client) |
| 111 | + const rng = mulberry32(plan.queries.length) |
| 112 | + const queries = plan.queries.map(render) |
| 113 | + const references = [] |
| 114 | + for (const q of queries) references.push(await outcome(clients.reference.query(q.text, q.values))) |
| 115 | + const full = references.map(({ result, error }) => summarize(result, error)) |
| 116 | + |
| 117 | + // the whole round in flight at once on the pipelined connection |
| 118 | + const pipelined = await Promise.all(queries.map((q) => settle(clients.pipeline.query(q.text, q.values)))) |
| 119 | + for (let i = 0; i < queries.length; i++) { |
| 120 | + if (pipelined[i] !== full[i]) |
| 121 | + return `pipeline, query ${i} |
| 122 | + reference: ${full[i]} |
| 123 | + pipeline: ${pipelined[i]}` |
| 124 | + } |
| 125 | + |
| 126 | + for (const [arm, runArm] of Object.entries(ARMS)) { |
| 127 | + for (let i = 0; i < queries.length; i++) { |
| 128 | + const got = await runArm(clients[arm], queries[i], rng, references[i]) |
| 129 | + const { result, error } = references[i] |
| 130 | + const expected = arm === 'cursor' || arm === 'stream' ? rowsOnly(result, error) : full[i] |
| 131 | + if (got !== expected) |
| 132 | + return `${arm}, query ${i} |
| 133 | + reference: ${expected} |
| 134 | + ${arm}: ${got}` |
| 135 | + } |
| 136 | + } |
| 137 | + return null |
| 138 | +} |
| 139 | + |
| 140 | +const close = async () => { |
| 141 | + for (const client of Object.values(clients)) await client.end() |
| 142 | +} |
| 143 | + |
| 144 | +main({ name: 'modes', draw, run, variants, source, close }) |
0 commit comments