perf(postgres): name the fixed kv statement shapes (#2173)

This commit is contained in:
Hampus
2026-08-30 23:01:21 +02:00
committed by GitHub
parent f0e7c25e4c
commit 8ff6518797
3 changed files with 354 additions and 9 deletions
+15 -3
View File
@@ -17,7 +17,11 @@ interface PostgresConfig {
}
export interface PostgresQueryable {
query<T extends QueryResultRow = QueryResultRow>(text: string, values?: Array<unknown>): Promise<QueryResult<T>>;
query<T extends QueryResultRow = QueryResultRow>(
text: string,
values?: Array<unknown>,
name?: string,
): Promise<QueryResult<T>>;
}
export interface IPostgresClient extends PostgresQueryable {
@@ -92,15 +96,16 @@ class PostgresClient implements IPostgresClient {
async query<T extends QueryResultRow = QueryResultRow>(
text: string,
values: Array<unknown> = [],
name?: string,
): Promise<QueryResult<T>> {
return this.getPool().query<T>(text, values);
return this.getPool().query<T>({text, values, name});
}
async transaction<T>(fn: (client: PostgresQueryable) => Promise<T>): Promise<T> {
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: <T extends QueryResultRow = QueryResultRow>(text: string, values: Array<unknown> = [], name?: string) =>
client.query<T>({text, values, name}),
};
}
async function rollback(client: PoolClient): Promise<void> {
try {
await client.query('ROLLBACK');
@@ -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<Array<StoredRow>> {
const result = await db.query<StoredRow>(
`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<string, StoredRow>();
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<Row | null> {
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;
@@ -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<string, unknown>;
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<void> {
await new Promise((resolve) => setTimeout(resolve, ms));
}
async function freePort(): Promise<number> {
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<Row> = {
name: 'stmt_composite',
columns: ['owner_id', 'item_id', 'payload'],
primaryKey: ['owner_id', 'item_id'],
partitionKey: ['owner_id'],
};
const Bucketed: KvTableSpec<Row> = {
name: 'stmt_bucketed',
columns: ['bucket', 'item_id', 'payload'],
primaryKey: ['item_id'],
partitionKey: ['bucket'],
};
function recordingClient(statements: Array<Statement>): IPostgresClient {
const client = {
async query(text: string, _values?: Array<unknown>, 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<unknown>) {
return fn(client);
},
kvTable() {
return KV_TABLE;
},
} as never;
return client;
}
function meta(spec: KvTableSpec<Row>, action: string, where: Array<WhereExpr<Row>>, extra: Row = {}): KvQueryMeta<Row> {
return {action, table: spec, where, columns: spec.columns, ...extra} as unknown as KvQueryMeta<Row>;
}
const eq = (col: string): WhereExpr<Row> => ({kind: 'eq', col, param: col}) as WhereExpr<Row>;
const isIn = (col: string): WhereExpr<Row> => ({kind: 'in', col, param: col}) as WhereExpr<Row>;
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<Array<Statement>> {
const statements: Array<Statement> = [];
const executor = new PostgresKvQueryExecutor(recordingClient(statements));
const cases: Array<[KvQueryMeta<Row>, 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<string, string>();
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<string, Set<string>>();
for (const statement of await runShapes()) {
if (statement.name === undefined) continue;
const names = namesByText.get(statement.text) ?? new Set<string>();
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<Row>({
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<Row>({
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<Row>({
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<Row>({
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<Row>({
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<Row>({
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<Row>({
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',
]);
});
});