From 8ff6518797ab301f33cff2a60a8b0ae25bd9d66c Mon Sep 17 00:00:00 2001 From: Hampus Date: Sun, 30 Aug 2026 23:01:21 +0200 Subject: [PATCH] perf(postgres): name the fixed kv statement shapes (#2173) --- fluxer_api/pkgs/postgres/src/Client.ts | 18 +- .../api/database/PostgresKvQueryExecutor.ts | 33 +- .../database/PostgresKvStatementNames.test.ts | 312 ++++++++++++++++++ 3 files changed, 354 insertions(+), 9 deletions(-) create mode 100644 fluxer_api/src/api/database/PostgresKvStatementNames.test.ts diff --git a/fluxer_api/pkgs/postgres/src/Client.ts b/fluxer_api/pkgs/postgres/src/Client.ts index ec822cd4a..d488da930 100644 --- a/fluxer_api/pkgs/postgres/src/Client.ts +++ b/fluxer_api/pkgs/postgres/src/Client.ts @@ -17,7 +17,11 @@ interface PostgresConfig { } export interface PostgresQueryable { - query(text: string, values?: Array): Promise>; + query( + text: string, + values?: Array, + name?: string, + ): Promise>; } export interface IPostgresClient extends PostgresQueryable { @@ -92,15 +96,16 @@ class PostgresClient implements IPostgresClient { async query( text: string, values: Array = [], + name?: string, ): Promise> { - return this.getPool().query(text, values); + return this.getPool().query({text, values, name}); } async transaction(fn: (client: PostgresQueryable) => Promise): Promise { const client = await this.getPool().connect(); try { await client.query('BEGIN'); - const result = await fn(client); + const result = await fn(poolClientQueryable(client)); await client.query('COMMIT'); return result; } catch (error) { @@ -123,6 +128,13 @@ class PostgresClient implements IPostgresClient { } } +function poolClientQueryable(client: PoolClient): PostgresQueryable { + return { + query: (text: string, values: Array = [], name?: string) => + client.query({text, values, name}), + }; +} + async function rollback(client: PoolClient): Promise { try { await client.query('ROLLBACK'); diff --git a/fluxer_api/src/api/database/PostgresKvQueryExecutor.ts b/fluxer_api/src/api/database/PostgresKvQueryExecutor.ts index d3f97b80a..4b9e8c0c5 100644 --- a/fluxer_api/src/api/database/PostgresKvQueryExecutor.ts +++ b/fluxer_api/src/api/database/PostgresKvQueryExecutor.ts @@ -67,6 +67,17 @@ const EXPIRED_STORED_ROW = 'kv.expires_at IS NOT NULL AND kv.expires_at <= now() const MERGED_ROW_DATA = `CASE WHEN ${EXPIRED_STORED_ROW} THEN EXCLUDED.row_data ELSE kv.row_data || EXCLUDED.row_data END`; const KEPT_EXPIRES_AT = `CASE WHEN ${EXPIRED_STORED_ROW} THEN NULL ELSE kv.expires_at END`; +function planStatementName(prefix: string, plan: CandidatePlan): string | undefined { + switch (plan.kind) { + case 'rowKeys': + return `${prefix}_rowkeys`; + case 'range': + return `${prefix}_range`; + default: + return undefined; + } +} + function normalizeCql(cql: string): string { return cql.replace(/\s+/g, ' ').trim(); } @@ -862,10 +873,12 @@ export class PostgresKvQueryExecutor { meta: KvQueryMeta, fragments: PlanFragments, db: PostgresQueryable, + name?: string, ): Promise> { const result = await db.query( `SELECT kv.row_key, kv.row_data FROM ${this.table} kv WHERE kv.table_name = $1${fragments.predicate} AND (kv.expires_at IS NULL OR kv.expires_at > now())`, [meta.table.name, ...fragments.params], + name, ); return result.rows; } @@ -874,10 +887,11 @@ export class PostgresKvQueryExecutor { if (plan.candidates.kind === 'none') return []; if (plan.candidates.kind === 'scan') logFullScan(meta); const groups = planFragmentGroups(plan.candidates); - if (groups.length === 1) return this.candidateGroup(meta, groups[0]!, db); + const name = planStatementName('kv_sel', plan.candidates); + if (groups.length === 1) return this.candidateGroup(meta, groups[0]!, db, name); const byRowKey = new Map(); for (const fragments of groups) { - for (const stored of await this.candidateGroup(meta, fragments, db)) byRowKey.set(stored.row_key, stored); + for (const stored of await this.candidateGroup(meta, fragments, db, name)) byRowKey.set(stored.row_key, stored); } return [...byRowKey.values()]; } @@ -913,6 +927,7 @@ export class PostgresKvQueryExecutor { const result = await db.query<{count: string}>( `SELECT count(*) AS count FROM ${this.table} kv WHERE kv.table_name = $1${fragments.predicate} AND (kv.expires_at IS NULL OR kv.expires_at > now())`, [meta.table.name, ...fragments.params], + planStatementName('kv_count', plan.candidates), ); total += Number(result.rows[0]?.count ?? 0); } @@ -929,6 +944,7 @@ export class PostgresKvQueryExecutor { await db.query( `DELETE FROM ${this.table} WHERE table_name = $1 AND row_key = $2 AND expires_at IS NOT NULL AND expires_at <= now()`, [meta.table.name, key], + 'kv_del_expired', ); } const expiresAt = ttlExpiresAt(meta, params) ?? null; @@ -946,6 +962,7 @@ WHERE NOT $6`, expiresAt, meta.ifNotExists === true, ], + 'kv_upsert', ); if (meta.ifNotExists) { return [{'[applied]': result.rowCount === 1}]; @@ -967,6 +984,7 @@ VALUES ($1, $2, $3, $4::jsonb, $5, now()) ON CONFLICT (table_name, row_key) DO UPDATE SET partition_key = EXCLUDED.partition_key, row_data = ${MERGED_ROW_DATA}, expires_at = ${expiresAtExpr}, updated_at = now()`, [meta.table.name, partitionKey(meta, incoming), key, JSON.stringify(encodeRow(incoming)), ttl ?? null], + ttl === undefined ? 'kv_patch_keep_ttl' : 'kv_patch_set_ttl', ); } @@ -979,6 +997,7 @@ DO UPDATE SET partition_key = EXCLUDED.partition_key, row_data = ${MERGED_ROW_DA await db.query( `DELETE FROM ${this.table} kv WHERE kv.table_name = $1${fragments.predicate} AND (kv.expires_at IS NULL OR kv.expires_at > now())`, [meta.table.name, ...fragments.params], + planStatementName('kv_del', plan.candidates), ); } return; @@ -995,16 +1014,18 @@ DO UPDATE SET partition_key = EXCLUDED.partition_key, row_data = ${MERGED_ROW_DA ) .map((stored) => stored.row_key); if (matchingKeys.length === 0) return; - await db.query(`DELETE FROM ${this.table} WHERE table_name = $1 AND row_key = ANY($2::text[])`, [ - meta.table.name, - matchingKeys, - ]); + await db.query( + `DELETE FROM ${this.table} WHERE table_name = $1 AND row_key = ANY($2::text[])`, + [meta.table.name, matchingKeys], + 'kv_del_keys', + ); } private async getRow(meta: KvQueryMeta, key: string, db: PostgresQueryable): Promise { const result = await db.query<{row_data: unknown}>( `SELECT row_data FROM ${this.table} WHERE table_name = $1 AND row_key = $2 AND (expires_at IS NULL OR expires_at > now()) LIMIT 1`, [meta.table.name, key], + 'kv_get_row', ); const row = result.rows[0]; return row ? decodeRow(row.row_data) : null; diff --git a/fluxer_api/src/api/database/PostgresKvStatementNames.test.ts b/fluxer_api/src/api/database/PostgresKvStatementNames.test.ts new file mode 100644 index 000000000..1f59f85e4 --- /dev/null +++ b/fluxer_api/src/api/database/PostgresKvStatementNames.test.ts @@ -0,0 +1,312 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +import {execFileSync, spawnSync} from 'node:child_process'; +import {createServer} from 'node:net'; +import { + getDefaultPostgresClient, + type IPostgresClient, + initPostgres, + shutdownPostgres, +} from '@pkgs/postgres/src/Client'; +import {afterAll, beforeAll, describe, expect, it} from 'vitest'; +import type {CassandraParams, KvQueryMeta, KvTableSpec, WhereExpr} from './CassandraTypes'; +import {ensurePostgresKvSchema, PostgresKvQueryExecutor} from './PostgresKvQueryExecutor'; + +type Row = Record; + +interface Statement { + text: string; + name: string | undefined; +} + +const KV_TABLE = 'kv_stmt_names'; +const CONTAINER = `fluxer-kvstmt-${process.pid.toString(36)}-${Date.now().toString(36)}`; +const dockerAvailable = spawnSync('docker', ['version'], {stdio: 'ignore'}).status === 0; + +async function sleep(ms: number): Promise { + await new Promise((resolve) => setTimeout(resolve, ms)); +} + +async function freePort(): Promise { + return new Promise((resolve, reject) => { + const server = createServer(); + server.on('error', reject); + server.listen(0, '127.0.0.1', () => { + const address = server.address(); + if (typeof address === 'string' || address === null) { + reject(new Error('no port')); + return; + } + const port = address.port; + server.close(() => resolve(port)); + }); + }); +} + +const Composite: KvTableSpec = { + name: 'stmt_composite', + columns: ['owner_id', 'item_id', 'payload'], + primaryKey: ['owner_id', 'item_id'], + partitionKey: ['owner_id'], +}; + +const Bucketed: KvTableSpec = { + name: 'stmt_bucketed', + columns: ['bucket', 'item_id', 'payload'], + primaryKey: ['item_id'], + partitionKey: ['bucket'], +}; + +function recordingClient(statements: Array): IPostgresClient { + const client = { + async query(text: string, _values?: Array, name?: string) { + statements.push({text, name}); + const rows = text.startsWith('SELECT kv.row_key') + ? [{row_key: 'stmt_row', row_data: {owner_id: 'o1', item_id: 'i1', payload: 'p1'}}] + : []; + return {rows, rowCount: rows.length}; + }, + async connect() {}, + async shutdown() {}, + isConnected() { + return true; + }, + async transaction(fn: (db: unknown) => Promise) { + return fn(client); + }, + kvTable() { + return KV_TABLE; + }, + } as never; + return client; +} + +function meta(spec: KvTableSpec, action: string, where: Array>, extra: Row = {}): KvQueryMeta { + return {action, table: spec, where, columns: spec.columns, ...extra} as unknown as KvQueryMeta; +} + +const eq = (col: string): WhereExpr => ({kind: 'eq', col, param: col}) as WhereExpr; +const isIn = (col: string): WhereExpr => ({kind: 'in', col, param: col}) as WhereExpr; + +const OWNER_ITEM = {owner_id: 'o1', item_id: 'i1', payload: 'p1'} as CassandraParams; +const BUCKET = {bucket: 'b1', item_id: 'i1', payload: 'p1'} as CassandraParams; + +async function runShapes(): Promise> { + const statements: Array = []; + const executor = new PostgresKvQueryExecutor(recordingClient(statements)); + const cases: Array<[KvQueryMeta, CassandraParams]> = [ + [meta(Composite, 'select', [eq('owner_id'), eq('item_id')]), OWNER_ITEM], + [meta(Composite, 'select', [eq('owner_id')]), OWNER_ITEM], + [meta(Composite, 'select', [isIn('owner_id')]), {owner_id: ['o1', 'o2']} as CassandraParams], + [meta(Composite, 'select', []), {} as CassandraParams], + [meta(Bucketed, 'select', [eq('bucket')]), BUCKET], + [meta(Bucketed, 'select', [isIn('bucket')]), {bucket: ['b1', 'b2']} as CassandraParams], + [meta(Composite, 'count', [eq('owner_id'), eq('item_id')]), OWNER_ITEM], + [meta(Composite, 'count', [eq('owner_id')]), OWNER_ITEM], + [meta(Composite, 'count', [isIn('owner_id')]), {owner_id: ['o1', 'o2']} as CassandraParams], + [meta(Composite, 'count', []), {} as CassandraParams], + [meta(Bucketed, 'count', [eq('bucket')]), BUCKET], + [meta(Bucketed, 'count', [isIn('bucket')]), {bucket: ['b1', 'b2']} as CassandraParams], + [meta(Composite, 'delete', [eq('owner_id'), eq('item_id')]), OWNER_ITEM], + [meta(Composite, 'delete', [eq('owner_id'), eq('payload')]), OWNER_ITEM], + [meta(Composite, 'upsert', []), OWNER_ITEM], + [meta(Composite, 'upsert', [], {ifNotExists: true}), OWNER_ITEM], + [meta(Composite, 'patch', [eq('owner_id'), eq('item_id')], {patchKeys: ['payload']}), OWNER_ITEM], + [ + meta(Composite, 'patch', [eq('owner_id'), eq('item_id')], {patchKeys: ['payload'], ttlParamName: 'ttl_'}), + {...OWNER_ITEM, ttl_: 600} as CassandraParams, + ], + ]; + for (const [kvMeta, params] of cases) { + await executor.executeQuery({cql: `__stmt_${kvMeta.action}`, params, kvMeta: kvMeta as KvQueryMeta}); + } + return statements; +} + +describe('PostgresKvQueryExecutor statement names', () => { + it('names every key-pinned statement shape exactly once', async () => { + const statements = await runShapes(); + const named = new Map(); + for (const statement of statements) { + if (statement.name === undefined) continue; + const seen = named.get(statement.name); + expect(seen ?? statement.text).toBe(statement.text); + named.set(statement.name, statement.text); + } + expect([...named.keys()].sort()).toEqual([ + 'kv_count_range', + 'kv_count_rowkeys', + 'kv_del_expired', + 'kv_del_keys', + 'kv_del_rowkeys', + 'kv_get_row', + 'kv_patch_keep_ttl', + 'kv_patch_set_ttl', + 'kv_sel_range', + 'kv_sel_rowkeys', + 'kv_upsert', + ]); + }); + + it('leaves the OR-of-ranges, partition and scan shapes unnamed', async () => { + const statements = await runShapes(); + const shapes = [' OR (', 'kv.partition_key = $2', 'kv.partition_key = ANY($2::text[])', '$1 AND (kv.expires_at']; + for (const shape of shapes) { + const matching = statements.filter((statement) => statement.text.includes(shape)); + expect(matching.length).toBeGreaterThan(0); + expect(matching.filter((statement) => statement.name !== undefined)).toEqual([]); + } + }); + + it('never reuses one statement text under two names', async () => { + const namesByText = new Map>(); + for (const statement of await runShapes()) { + if (statement.name === undefined) continue; + const names = namesByText.get(statement.text) ?? new Set(); + names.add(statement.name); + namesByText.set(statement.text, names); + } + expect([...namesByText.values()].filter((names) => names.size > 1)).toEqual([]); + }); +}); + +describe.skipIf(!dockerAvailable)('PostgresKvQueryExecutor statement names against postgres', () => { + let raw: IPostgresClient; + let executor: PostgresKvQueryExecutor; + + beforeAll(async () => { + const port = await freePort(); + execFileSync( + 'docker', + [ + 'run', + '-d', + '--name', + CONTAINER, + '-e', + 'POSTGRES_USER=fluxer', + '-e', + 'POSTGRES_PASSWORD=fluxer', + '-e', + 'POSTGRES_DB=fluxer', + '-p', + `127.0.0.1:${port}:5432`, + 'postgres:16-alpine', + '-c', + 'fsync=off', + ], + {stdio: 'ignore'}, + ); + let ready = false; + for (let attempt = 0; attempt < 180 && !ready; attempt += 1) { + await sleep(500); + const probe = spawnSync('docker', ['exec', CONTAINER, 'pg_isready', '-U', 'fluxer', '-d', 'fluxer'], { + stdio: 'ignore', + }); + if (probe.status !== 0) continue; + try { + await initPostgres({ + url: `postgres://fluxer:fluxer@127.0.0.1:${port}/fluxer`, + maxConnections: 1, + kvTable: KV_TABLE, + }); + await getDefaultPostgresClient().query('SELECT 1'); + ready = true; + } catch { + await shutdownPostgres().catch(() => {}); + } + } + if (!ready) throw new Error('postgres never came up'); + raw = getDefaultPostgresClient(); + await ensurePostgresKvSchema(raw); + executor = new PostgresKvQueryExecutor(raw); + }, 900_000); + + afterAll(async () => { + await shutdownPostgres().catch(() => {}); + spawnSync('docker', ['rm', '-f', CONTAINER], {stdio: 'ignore'}); + }); + + it('prepares each named shape server side and keeps reading the right rows', async () => { + for (let index = 0; index < 4; index += 1) { + await executor.executeQuery({ + cql: '__stmt_seed', + params: {owner_id: `o${index % 2}`, item_id: `i${index}`, payload: `p${index}`} as CassandraParams, + kvMeta: meta(Composite, 'upsert', []) as KvQueryMeta, + }); + } + for (let index = 0; index < 12; index += 1) { + const point = await executor.executeQuery({ + cql: '__stmt_point', + params: {owner_id: 'o0', item_id: 'i0'} as CassandraParams, + kvMeta: meta(Composite, 'select', [eq('owner_id'), eq('item_id')]) as KvQueryMeta, + }); + expect(point.map((row) => row.payload)).toEqual(['p0']); + const range = await executor.executeQuery({ + cql: '__stmt_range', + params: {owner_id: 'o0'} as CassandraParams, + kvMeta: meta(Composite, 'select', [eq('owner_id')]) as KvQueryMeta, + }); + expect(range).toHaveLength(2); + } + await executor.executeQuery({ + cql: '__stmt_patch', + params: {owner_id: 'o0', item_id: 'i0', payload: 'patched'} as CassandraParams, + kvMeta: meta(Composite, 'patch', [eq('owner_id'), eq('item_id')], {patchKeys: ['payload']}) as KvQueryMeta, + }); + const reapplied = await executor.executeQuery({ + cql: '__stmt_lwt', + params: {owner_id: 'o0', item_id: 'i0', payload: 'ignored'} as CassandraParams, + kvMeta: meta(Composite, 'upsert', [], {ifNotExists: true}) as KvQueryMeta, + }); + expect(reapplied).toEqual([{'[applied]': false}]); + const claimed = await executor.executeQuery({ + cql: '__stmt_lwt', + params: {owner_id: 'o9', item_id: 'i9', payload: 'claimed'} as CassandraParams, + kvMeta: meta(Composite, 'upsert', [], {ifNotExists: true}) as KvQueryMeta, + }); + expect(claimed).toEqual([{'[applied]': true}]); + await executor.executeQuery({ + cql: '__stmt_patch_ttl', + params: {owner_id: 'o9', item_id: 'i9', payload: 'expiring', ttl_: 600} as CassandraParams, + kvMeta: meta(Composite, 'patch', [eq('owner_id'), eq('item_id')], { + patchKeys: ['payload'], + ttlParamName: 'ttl_', + }) as KvQueryMeta, + }); + const expiring = await executor.executeQuery({ + cql: '__stmt_point', + params: {owner_id: 'o9', item_id: 'i9'} as CassandraParams, + kvMeta: meta(Composite, 'select', [eq('owner_id'), eq('item_id')]) as KvQueryMeta, + }); + expect(expiring.map((row) => row.payload)).toEqual(['expiring']); + const patched = await executor.executeQuery({ + cql: '__stmt_point', + params: {owner_id: 'o0', item_id: 'i0'} as CassandraParams, + kvMeta: meta(Composite, 'select', [eq('owner_id'), eq('item_id')]) as KvQueryMeta, + }); + expect(patched.map((row) => row.payload)).toEqual(['patched']); + await executor.executeQuery({ + cql: '__stmt_delete', + params: {owner_id: 'o0', item_id: 'i0'} as CassandraParams, + kvMeta: meta(Composite, 'delete', [eq('owner_id'), eq('item_id')]) as KvQueryMeta, + }); + const remaining = await executor.executeQuery({ + cql: '__stmt_range', + params: {owner_id: 'o0'} as CassandraParams, + kvMeta: meta(Composite, 'select', [eq('owner_id')]) as KvQueryMeta, + }); + expect(remaining).toHaveLength(1); + const prepared = await raw.query<{name: string}>('SELECT name FROM pg_prepared_statements ORDER BY name'); + expect(prepared.rows.map((row) => row.name)).toEqual([ + 'kv_del_expired', + 'kv_del_rowkeys', + 'kv_get_row', + 'kv_patch_keep_ttl', + 'kv_patch_set_ttl', + 'kv_sel_range', + 'kv_sel_rowkeys', + 'kv_upsert', + ]); + }); +});