import pg from 'pg'; import type { AppConfig } from './config.js'; import { writeEmergencyDiagnostic } from './diagnostics.js'; const { Pool } = pg; let pool: pg.Pool | null = null; export const REQUIRED_SCHEMA_VERSION = '021_admin_task_data_plane_isolation'; export interface DatabaseReadiness { ready: boolean; database: boolean; schema: boolean; requiredMigration: string; errorCode?: string; } export function quoteIdentifier(identifier: string): string { return `"${identifier.replace(/"/g, '""')}"`; } function logDatabaseError(context: string, error: unknown): void { const event = context === 'idle pool client error' ? 'database.pool.idle_client_error' : 'database.transaction.rollback_failed'; writeEmergencyDiagnostic(event, error, { diagnostic_stage: 'database' }); } export class DatabaseStartupError extends Error { constructor(public readonly code: string, message: string) { super(message); this.name = 'DatabaseStartupError'; } } export function getPool(config: AppConfig): pg.Pool { if (!pool) { pool = new Pool({ connectionString: config.DATABASE_URL, max: 10, idleTimeoutMillis: 30_000, connectionTimeoutMillis: 5_000, application_name: 'liansyn-platform', options: `-c search_path=${quoteIdentifier(config.DATABASE_SCHEMA)},public`, ...(config.DATABASE_SSL === 'require' ? { ssl: { rejectUnauthorized: false } } : {}) }); // pg removes an idle client that emits an error, but the Pool itself must // have an error listener. Without one, a transient network timeout on an // idle remote connection is emitted as an unhandled Node error and takes // down the whole control plane. pool.on('error', (error) => logDatabaseError('idle pool client error', error)); } return pool; } export async function withTransaction(config: AppConfig, fn: (client: pg.PoolClient) => Promise): Promise { const client = await getPool(config).connect(); try { await client.query('BEGIN'); const result = await fn(client); await client.query('COMMIT'); return result; } catch (error) { try { await client.query('ROLLBACK'); } catch (rollbackError) { logDatabaseError('transaction rollback failed', rollbackError); } throw error; } finally { client.release(); } } export async function closePool(): Promise { if (!pool) return; const active = pool; pool = null; await active.end(); } export async function databaseReady(config: AppConfig): Promise { const readiness = await databaseReadiness(config); return readiness.ready; } export async function databaseReadiness(config: AppConfig): Promise { const readiness: DatabaseReadiness = { ready: false, database: false, schema: false, requiredMigration: REQUIRED_SCHEMA_VERSION }; try { await getPool(config).query('SELECT 1'); readiness.database = true; const result = await getPool(config).query( `SELECT EXISTS ( SELECT 1 FROM schema_migrations WHERE version = $1 ) AS applied`, [REQUIRED_SCHEMA_VERSION] ); readiness.schema = result.rows[0]?.applied === true; readiness.ready = readiness.database && readiness.schema; } catch (error) { const code = error && typeof error === 'object' && 'code' in error ? String((error as { code?: unknown }).code || '') : ''; if (code) readiness.errorCode = code; } return readiness; } export async function assertDatabaseSchema(config: AppConfig): Promise { const readiness = await databaseReadiness(config); if (!readiness.database) { throw new DatabaseStartupError( readiness.errorCode || 'database_unavailable', `数据库不可用,无法启动控制平面${readiness.errorCode ? `(${readiness.errorCode})` : ''}。` ); } if (!readiness.schema) { throw new DatabaseStartupError( 'database_schema_outdated', `数据库迁移 ${REQUIRED_SCHEMA_VERSION} 尚未应用,请先执行 npm run db:migrate。` ); } }