import pg from "pg"; import { AsyncLocalStorage } from "async_hooks"; import { drizzle, type NodePgDatabase } from "drizzle-orm/node-postgres"; import * as schema from "@shared/schema"; type Schema = typeof schema; let _pool: pg.Pool | null = null; let _globalDb: NodePgDatabase | null = null; type TenantCtx = { tenantDb: NodePgDatabase; client: pg.PoolClient }; // Per-request / per-operation tenant context. // Stores a Drizzle instance bound to a dedicated pg.PoolClient. // The `db` proxy routes all db.* calls to that instance when a context is active. export const _tenantCtx = new AsyncLocalStorage(); function getPool(): pg.Pool { if (_pool) return _pool; if (!process.env.DATABASE_URL) { const msg = "[DB] FATAL: DATABASE_URL is not set.\n" + " Make sure your .env file (or environment) defines DATABASE_URL.\n" + " Example: DATABASE_URL=postgres://user:password@host:5432/dbname"; console.error(msg); throw new Error(msg); } // SSL включается для sslmode=require и Neon. rejectUnauthorized управляется // env DATABASE_SSL_REJECT_UNAUTHORIZED (дефолт true — проверка сертификата); // 'false' допустимо только с самоподписанными сертификатами (dev-стенды). const sslRequired = process.env.DATABASE_URL.includes("sslmode=require") || process.env.DATABASE_URL.includes("neon.tech"); const sslRejectUnauthorized = process.env.DATABASE_SSL_REJECT_UNAUTHORIZED !== "false"; _pool = new pg.Pool({ connectionString: process.env.DATABASE_URL, max: 50, idleTimeoutMillis: 30_000, connectionTimeoutMillis: 5_000, ssl: sslRequired ? { rejectUnauthorized: sslRejectUnauthorized } : undefined, }); _pool.on("error", (err) => { let host = ""; try { host = new URL(process.env.DATABASE_URL!).host; } catch { /* ignore */ } console.error( `[DB] ERROR: Unexpected database connection error (host: ${host}).\n` + " The app may be unable to serve requests until the database is reachable.\n" + ` Details: ${err.message}` ); }); setInterval(() => { const p = getPool(); console.log(`[POOL] total=${p.totalCount} idle=${p.idleCount} waiting=${p.waitingCount}`); }, 60_000); const sigHandler = async () => { console.log("[SHUTDOWN] Draining pool..."); if (_pool) { await _pool.end(); } process.exit(0); }; process.on("SIGTERM", sigHandler); return _pool; } export const pool: pg.Pool = new Proxy({} as pg.Pool, { get(_target, prop) { const p = getPool(); const value = Reflect.get(p, prop, p); return typeof value === "function" ? value.bind(p) : value; }, }); // Context-aware Drizzle instance. // When inside an openTenantCtx/openSuperAdminCtx context, routes all db.* calls // to the dedicated client with SET SESSION already applied. // Falls back to the shared pool otherwise. export const db: NodePgDatabase = new Proxy({} as NodePgDatabase, { get(_target, prop) { const ctx = _tenantCtx.getStore(); const target = ctx ? ctx.tenantDb : (() => { if (!_globalDb) _globalDb = drizzle({ client: getPool(), schema }); return _globalDb; })(); const value = Reflect.get(target, prop, target); return typeof value === "function" ? value.bind(target) : value; }, }); // Handle returned by openTenantCtx / openSuperAdminCtx. export type CtxHandle = { run: (fn: () => void) => void; release: () => Promise; }; // Opens a dedicated client scoped to one organization using SET SESSION. // SET SESSION (not SET LOCAL) persists the variable across transaction boundaries, // so storage methods calling db.transaction() do not lose the RLS context. // The client is released (with session var reset) when release() is called. export async function openTenantCtx(organizationId: number): Promise { const client = await getPool().connect(); let released = false; try { // set_config(name, value, is_local): false = session scope (survives db.transaction()) await client.query("SELECT set_config('app.current_org_id', $1, false)", [String(organizationId)]); const tenantDb = drizzle({ client, schema }); const ctx: TenantCtx = { tenantDb, client }; return { run: (fn) => _tenantCtx.run(ctx, fn), release: async () => { if (released) return; released = true; try { await client.query("SELECT set_config('app.current_org_id', '', false)"); } catch { /* ignore */ } client.release(); }, }; } catch (err) { client.release(); throw err; } } // Opens a dedicated client with superadmin bypass using SET SESSION. // SET SESSION persists is_superadmin across any db.transaction() calls in storage. export async function openSuperAdminCtx(): Promise { const client = await getPool().connect(); let released = false; try { await client.query("SELECT set_config('app.is_superadmin', 'true', false)"); const saDb = drizzle({ client, schema }); const ctx: TenantCtx = { tenantDb: saDb, client }; return { run: (fn) => _tenantCtx.run(ctx, fn), release: async () => { if (released) return; released = true; try { await client.query("SELECT set_config('app.is_superadmin', 'false', false)"); } catch { /* ignore */ } client.release(); }, }; } catch (err) { client.release(); throw err; } } // Executes fn within a dedicated client scoped to one organization. // Uses BEGIN + set_config(is_local=true) for short-lived point operations. // Re-entrant: reuses the existing context if already inside one. Note: if called // inside a withSuperAdmin context, the superadmin scope is reused (no narrowing). export async function withTenant( organizationId: number, fn: (tenantDb: NodePgDatabase) => Promise ): Promise { const existing = _tenantCtx.getStore(); if (existing) return fn(existing.tenantDb); const client = await getPool().connect(); try { await client.query("BEGIN"); // set_config(name, value, is_local=true): transaction-local scope, reset on ROLLBACK/COMMIT await client.query("SELECT set_config('app.current_org_id', $1, true)", [String(organizationId)]); const tenantDb = drizzle({ client, schema }); const result = await _tenantCtx.run({ tenantDb, client }, () => fn(tenantDb)); await client.query("COMMIT"); return result; } catch (err) { try { await client.query("ROLLBACK"); } catch { /* ignore */ } throw err; } finally { client.release(); } } // Executes fn with superadmin bypass for cross-tenant data access. // Uses BEGIN + set_config(is_local=true) for short-lived point operations. // Re-entrant: reuses the existing context if already inside one. Note: if called // inside a withTenant context, the tenant scope is reused (no elevation). // Always call from explicitly superadmin-only code paths. export async function withSuperAdmin( fn: (saDb: NodePgDatabase) => Promise ): Promise { const existing = _tenantCtx.getStore(); if (existing) return fn(existing.tenantDb); const client = await getPool().connect(); try { await client.query("BEGIN"); await client.query("SELECT set_config('app.is_superadmin', 'true', true)"); const saDb = drizzle({ client, schema }); const result = await _tenantCtx.run({ tenantDb: saDb, client }, () => fn(saDb)); await client.query("COMMIT"); return result; } catch (err) { try { await client.query("ROLLBACK"); } catch { /* ignore */ } throw err; } finally { client.release(); } }