Files
CashuMints.space/api/src/db-postgres.ts
T
michilis 6f17b572b1 Expand ecash explorer capabilities
Add Fedimint discovery, dual SQLite/Postgres storage, richer review handling, and generated social imagery.
2026-08-21 02:10:48 +02:00

148 lines
4.5 KiB
TypeScript

import pg from 'pg';
import type { Db, Param, RunResult, Sql } from './db-driver.ts';
const { Pool, types } = pg;
/**
* node-postgres hands back BIGINT and NUMERIC as strings, because either can hold a
* value JavaScript's number cannot. Nothing here can: the counters are row counts, the
* timestamps are unix seconds, and the averages are ratings between 1 and 5. Left
* as-is the strings would flow straight into the JSON the API serves, turning
* `review_count: 12` into `"12"` and quietly breaking every comparison on the way.
*/
types.setTypeParser(20, (v) => Number.parseInt(v, 10)); // int8
types.setTypeParser(1700, (v) => Number.parseFloat(v)); // numeric
/**
* Rewrite `?` placeholders to `$1..$n`.
*
* Quoted text is copied verbatim so a `?` inside a string literal is not renumbered.
* Postgres escapes a quote by doubling it, which needs no special case: the closing
* quote ends one literal and the next character opens another.
*/
export function toPgPlaceholders(sql: string): string {
let out = '';
let quote: string | null = null;
let n = 0;
for (const ch of sql) {
if (quote !== null) {
if (ch === quote) quote = null;
out += ch;
} else if (ch === "'" || ch === '"') {
quote = ch;
out += ch;
} else if (ch === '?') {
out += `$${++n}`;
} else {
out += ch;
}
}
return out;
}
const translated = new Map<string, string>();
function translate(sql: string): string {
let pgSql = translated.get(sql);
if (pgSql === undefined) {
pgSql = toPgPlaceholders(sql);
translated.set(sql, pgSql);
}
return pgSql;
}
/** Bind a Sql surface to anything that can run a query: the pool, or one client. */
function surface(run: (text: string, values: Param[]) => Promise<pg.QueryResult>): Sql {
return {
async all<T>(sql: string, ...params: Param[]): Promise<T[]> {
return (await run(translate(sql), params)).rows as T[];
},
async get<T>(sql: string, ...params: Param[]): Promise<T | undefined> {
return (await run(translate(sql), params)).rows[0] as T | undefined;
},
async run(sql: string, ...params: Param[]): Promise<RunResult> {
return { changes: (await run(translate(sql), params)).rowCount ?? 0 };
},
};
}
/** Strip the password so a connection string is safe to put in a log line. */
export function redact(url: string): string {
try {
const parsed = new URL(url);
if (parsed.password) parsed.password = '***';
return parsed.toString();
} catch {
return 'postgres';
}
}
export class PostgresDb implements Db {
readonly dialect = 'postgres' as const;
readonly label: string;
#pool: pg.Pool;
#sql: Sql;
constructor(connectionString: string, poolMax: number) {
this.#pool = new Pool({ connectionString, max: poolMax });
// A pool client can die between checkouts — a Postgres restart, a killed session, an
// idle timeout on a proxy. Without a listener node-postgres emits that on the pool as
// an unhandled 'error' and takes the process down with it.
this.#pool.on('error', () => undefined);
this.#sql = surface((text, values) => this.#pool.query(text, values));
this.label = redact(connectionString);
}
all<T>(sql: string, ...params: Param[]): Promise<T[]> {
return this.#sql.all<T>(sql, ...params);
}
get<T>(sql: string, ...params: Param[]): Promise<T | undefined> {
return this.#sql.get<T>(sql, ...params);
}
run(sql: string, ...params: Param[]): Promise<RunResult> {
return this.#sql.run(sql, ...params);
}
async exec(statements: readonly string[]): Promise<void> {
// One connection for the lot: concurrent CREATE ... IF NOT EXISTS on the same
// catalog rows deadlocks in Postgres, and this runs on every boot.
const client = await this.#pool.connect();
try {
for (const sql of statements) await client.query(sql);
} finally {
client.release();
}
}
async transaction<T>(fn: (tx: Sql) => Promise<T>): Promise<T> {
const client = await this.#pool.connect();
const tx = surface((text, values) => client.query(text, values));
try {
await client.query('BEGIN');
const result = await fn(tx);
await client.query('COMMIT');
return result;
} catch (err) {
try {
await client.query('ROLLBACK');
} catch {
// The connection is already unusable; release() below discards it.
}
throw err;
} finally {
client.release();
}
}
close(): Promise<void> {
return this.#pool.end();
}
}