api-database-postgresql
Direct PostgreSQL access with node-postgres (pg) -- connection pools, parameterized queries, transactions, streaming, LISTEN/NOTIFY, error handling
Install
npx skills add https://github.com/agents-inc/skills/tree/main/dist/plugins/api-database-postgresql/skills/api-database-postgresql
claude plugin marketplace add https://llmmart.ai/marketplace.json && claude plugin install agents-inc-skills@llmmart
git clone https://github.com/agents-inc/skills.git
The skills CLI installs just this skill, for any of its supported agents. Claude Code installs the whole agents-inc/skills collection as a plugin from our marketplace. Git is the plain clone.
Skill manifest
PostgreSQL Patterns (node-postgres)
Quick Guide: Use the
pgpackage (v8.x) for direct PostgreSQL access. Always usePool-- never create individualClientinstances in application code. Use parameterized queries ($1,$2) for ALL user input -- never interpolate strings into SQL. For transactions, check out a dedicated client withpool.connect()and useBEGIN/COMMIT/ROLLBACKin atry/catch/finallythat always callsclient.release(). Handle the poolerrorevent to prevent process crashes from idle client errors. Usepg-query-streamfor large result sets to avoid loading everything into memory.
<critical_requirements>
CRITICAL: Before Using This Skill
All code must follow project conventions in CLAUDE.md (kebab-case, named exports, import ordering,
import type, named constants)
(You MUST use parameterized queries ($1, $2, ...) for ALL values -- NEVER concatenate or interpolate user input into SQL strings)
(You MUST use Pool for all database access -- NEVER create standalone Client instances in application code)
(You MUST release clients back to the pool in a finally block after pool.connect() -- leaked clients exhaust the pool and hang the application)
(You MUST handle the error event on Pool instances -- unhandled idle client errors crash the Node.js process)
</critical_requirements>
Examples
- Core Patterns -- Pool setup, parameterized queries, type-safe results, error handling
- Transactions -- BEGIN/COMMIT/ROLLBACK, savepoints, retry logic, advisory locks
- Streaming -- Cursors, pg-query-stream, LISTEN/NOTIFY for real-time
- Advanced -- SSL/TLS, prepared statements, migrations, testing patterns
Additional resources:
- reference.md -- Pool options, error codes, QueryResult properties, production checklist
Auto-detection: PostgreSQL, pg, node-postgres, Pool, Client, pool.query, pool.connect, client.query, $1, parameterized query, BEGIN, COMMIT, ROLLBACK, LISTEN, NOTIFY, pg_notify, pg-query-stream, pg-cursor, Cursor, QueryResult, QueryResultRow, connectionString, PGHOST, PGDATABASE, unique_violation, 23505, deadlock, 40P01, advisory lock
When to use:
- Direct SQL queries against PostgreSQL (not behind an ORM)
- Connection pool management for Node.js/PostgreSQL applications
- Transactions spanning multiple queries that must be atomic
- Streaming large result sets without loading everything into memory
- Real-time change notifications via LISTEN/NOTIFY
- Integration testing with transaction rollback isolation
Key patterns covered:
- Pool configuration and lifecycle (creation, error handling, graceful shutdown)
- Parameterized queries with
$1-style placeholders (SQL injection prevention) - Type-safe query results using TypeScript generics
- Transaction management with dedicated clients
- Streaming with pg-cursor and pg-query-stream
- LISTEN/NOTIFY for real-time PostgreSQL event handling
- PostgreSQL error code handling (constraint violations, deadlocks, serialization failures)
- SSL/TLS connection configuration
- Testing with transaction rollback isolation
When NOT to use:
- You need an ORM or query builder -- use your ORM/query builder skill instead
- You need in-memory caching -- use a caching solution
- You need document storage without relational constraints -- use a document database
- Simple key-value lookups at sub-millisecond latency -- use an in-memory data store
<decision_framework>
Decision Framework
pool.query() vs pool.connect()
Do I need a dedicated client?
├─ Single query, no transaction? -> pool.query() (auto-releases)
├─ Multiple queries in a transaction? -> pool.connect() + BEGIN/COMMIT/ROLLBACK
├─ LISTEN for notifications? -> pool.connect() (keep client for lifetime of listener)
├─ Cursor/streaming? -> pool.connect() (cursor binds to a connection)
└─ Prepared statements across queries? -> pool.connect() (plan caches per connection)
Error Handling Strategy
What kind of PostgreSQL error?
├─ 23505 (unique_violation)? -> Map to 409 Conflict, include constraint name
├─ 23503 (foreign_key_violation)? -> Map to 400 Bad Request, entity not found
├─ 23502 (not_null_violation)? -> Map to 400 Bad Request, missing required field
├─ 23514 (check_violation)? -> Map to 400 Bad Request, validation failed
├─ 40P01 (deadlock_detected)? -> Retry with backoff (safe to retry)
├─ 40001 (serialization_failure)? -> Retry with backoff (safe to retry)
├─ 57014 (query_canceled)? -> Timeout, consider increasing statement_timeout
├─ 08xxx (connection_exception)? -> Pool handles reconnection, log and retry
└─ Other? -> Log full error, return 500
Streaming Decision
How many rows will the query return?
├─ < 1,000 rows? -> pool.query() is fine (result fits in memory)
├─ 1,000 - 100,000 rows? -> pg-cursor with batch processing
├─ 100,000+ rows? -> pg-query-stream piped to a writable stream
└─ Need to export to file? -> pg-query-stream piped to file write stream
</decision_framework>
<red_flags>
RED FLAGS
High Priority Issues:
- Using string interpolation/concatenation for SQL values -- this is SQL injection, the most dangerous vulnerability in database code
- Using
pool.query()for transactions -- each call may use a different connection, so BEGIN/COMMIT have no effect - Not releasing clients after
pool.connect()-- leaked clients exhaust the pool; oncemaxclients leak, the app deadlocks onpool.connect() - Missing
pool.on("error")handler -- idle client errors are emitted on the pool; unhandled, they crash the Node.js process - Using standalone
Clientin application code -- no pooling, no reconnection, no concurrency
Medium Priority Issues:
SELECT *in production queries -- returns unnecessary columns, breaks when schema changes, prevents index-only scans- Loading millions of rows with
pool.query()instead of streaming -- causes memory exhaustion and GC pressure - Hardcoded connection strings -- prevents environment-specific configuration, risks credential leaks in version control
- Not handling specific PostgreSQL error codes -- generic error handling loses valuable information (which constraint, which column)
- Using
LISTENwithpool.query()-- notifications bind to a specific connection; pool.query releases the connection immediately
Common Mistakes:
- Forgetting that
result.rows[0]can beundefinedwhen no rows match -- always check before accessing - Relying on
result.rowCountforSELECTemptiness checks -- useresult.rows.lengthinstead;rowCountisnullfor some commands (e.g.,LOCK) androws.lengthis universally reliable - Using
$1inside string literals in SQL --'$1'is a literal string, not a parameter; use$1outside quotes - Forgetting that PostgreSQL arrays in parameters are automatically converted --
[1, 2, 3]becomes{1,2,3}which works for= ANY($1)but not forIN ($1)(use= ANY($1::int[])instead ofIN) - Calling
client.release(true)routinely -- passingtruedestroys the client instead of returning it to the pool; only use after unrecoverable errors
Gotchas & Edge Cases:
- Pool
errorevent vs query errors: Poolerrorfires for idle client backend disconnections (e.g., server restart). Query errors are thrown/rejected from the query call itself. You need both handlers. connectionTimeoutMillis: 0(default) means no timeout -- connections wait forever if the pool is exhausted. Always set a timeout in production.idleTimeoutMillisonly affects clients that have been returned to the pool -- a checked-out client that is never released will never be cleaned up.- PostgreSQL
numeric/decimaltypes are returned as strings by default (to avoid JavaScript floating-point precision loss). Parse them explicitly if you need numbers. LISTENsurvives transactions -- if youBEGIN,LISTEN channel,ROLLBACK, the listener is still active. LISTEN is not transactional.pool.end()waits for all checked-out clients to be released. If a client is leaked (never released),pool.end()hangs forever.- SSL connections: if the connection string contains any SSL parameters (
sslmode,sslcert,sslkey,sslrootcert), the entiresslconfig object is replaced -- use one or the other, not both.
</red_flags>
<critical_reminders>
CRITICAL REMINDERS
All code must follow project conventions in CLAUDE.md (kebab-case, named exports, import ordering,
import type, named constants)
(You MUST use parameterized queries ($1, $2, ...) for ALL values -- NEVER concatenate or interpolate user input into SQL strings)
(You MUST use Pool for all database access -- NEVER create standalone Client instances in application code)
(You MUST release clients back to the pool in a finally block after pool.connect() -- leaked clients exhaust the pool and hang the application)
(You MUST handle the error event on Pool instances -- unhandled idle client errors crash the Node.js process)
Failure to follow these rules will cause SQL injection vulnerabilities, connection pool exhaustion, application hangs, and process crashes.
</critical_reminders>
Files (skills)
-
examples
-
advanced.md 9.6 KB
# PostgreSQL -- Advanced Examples > Prepared statements, SSL/TLS configuration, migration patterns, and testing with transaction rollback. Reference from [SKILL.md](../SKILL.md). **Related examples:** - [core.md](core.md) -- Pool setup, parameterized queries, error handling - [transactions.md](transactions.md) -- Transactions, savepoints, retry logic - [streaming.md](streaming.md) -- Cursors, pg-query-stream, LISTEN/NOTIFY --- ## Prepared Statements Named queries cache execution plans on the PostgreSQL server per connection. Useful for complex queries executed repeatedly. ```typescript import type pg from "pg"; interface ProductRow { id: number; name: string; price: string; // numeric returns as string category: string; } const PRODUCT_SEARCH_QUERY = { name: "search-products", text: ` SELECT id, name, price, category FROM products WHERE category = $1 AND price BETWEEN $2 AND $3 ORDER BY price ASC LIMIT $4 `, }; const DEFAULT_PAGE_SIZE = 50; async function searchProducts( pool: pg.Pool, category: string, minPrice: number, maxPrice: number, limit: number = DEFAULT_PAGE_SIZE, ): Promise<ProductRow[]> { const result = await pool.query<ProductRow>({ ...PRODUCT_SEARCH_QUERY, values: [category, minPrice, maxPrice, limit], }); return result.rows; } export { searchProducts }; ``` **Why good:** Named query (`name: "search-products"`) tells PostgreSQL to cache the execution plan, subsequent calls skip parsing and planning. The query object is defined once as a constant. **When to use:** Complex queries with multiple JOINs where planning time is significant. For simple queries, the overhead of plan caching is not worth it. **Gotcha:** Prepared statement plans are cached per connection. If a pool rotates connections (due to `maxUses` or `maxLifetimeSeconds`), plans are lost on the new connection. This is usually fine -- the first execution on a new connection parses, subsequent ones are cached. **Gotcha:** If you run DDL (ALTER TABLE, CREATE INDEX) that changes the table structure, cached plans may become invalid. PostgreSQL automatically invalidates them, but be aware this causes a one-time re-plan. --- ## SSL/TLS Connections ```typescript import pg from "pg"; import { readFileSync } from "node:fs"; // Cloud-managed databases (AWS RDS, GCP Cloud SQL, etc.) // Most cloud providers use trusted CAs -- just enable SSL function createCloudPool(): pg.Pool { return new pg.Pool({ connectionString: process.env.DATABASE_URL, ssl: { rejectUnauthorized: true, // Verify server certificate }, }); } // Self-signed certificates (on-premise, custom CA) function createMtlsPool(): pg.Pool { return new pg.Pool({ connectionString: process.env.DATABASE_URL, ssl: { ca: readFileSync(process.env.PG_CA_CERT!), key: readFileSync(process.env.PG_CLIENT_KEY!), cert: readFileSync(process.env.PG_CLIENT_CERT!), rejectUnauthorized: true, }, }); } export { createCloudPool, createMtlsPool }; ``` **Why good:** `rejectUnauthorized: true` validates the server certificate (prevents MITM), mTLS for mutual authentication, cert paths from environment variables **Gotcha:** If the connection string contains any SSL parameters (`sslmode`, `sslcert`, `sslkey`, `sslrootcert`), the entire `ssl` config object is replaced and any additional options are lost. Use one or the other. **Gotcha:** `ssl: true` relies on Node.js TLS defaults (which reject unauthorized certs), but does not explicitly set `rejectUnauthorized`. Always use the object form (`ssl: { rejectUnauthorized: true }`) to be explicit about intent. --- ## Dynamic Password Callback Useful for AWS IAM database authentication or secrets managers that rotate credentials. ```typescript import pg from "pg"; function createPoolWithDynamicPassword( getPassword: () => Promise<string>, ): pg.Pool { return new pg.Pool({ host: process.env.PG_HOST, port: parseInt(process.env.PG_PORT ?? "5432", 10), database: process.env.PG_DATABASE, user: process.env.PG_USER, password: getPassword, // Called on each new connection ssl: { rejectUnauthorized: true }, }); } export { createPoolWithDynamicPassword }; ``` **Why good:** Password callback is invoked each time a new connection is established, so rotated credentials are picked up automatically without restarting the pool. --- ## Migration Pattern (Manual SQL Files) A simple, dependency-free migration approach using sequential SQL files. ```typescript import type pg from "pg"; import { readFileSync, readdirSync } from "node:fs"; import { join } from "node:path"; const MIGRATIONS_TABLE = "schema_migrations"; async function ensureMigrationsTable(pool: pg.Pool): Promise<void> { await pool.query(` CREATE TABLE IF NOT EXISTS ${MIGRATIONS_TABLE} ( id SERIAL PRIMARY KEY, name TEXT NOT NULL UNIQUE, applied_at TIMESTAMPTZ NOT NULL DEFAULT NOW() ) `); } async function getAppliedMigrations(pool: pg.Pool): Promise<Set<string>> { const result = await pool.query<{ name: string }>( `SELECT name FROM ${MIGRATIONS_TABLE} ORDER BY id`, ); return new Set(result.rows.map((r) => r.name)); } async function runMigrations( pool: pg.Pool, migrationsDir: string, ): Promise<string[]> { await ensureMigrationsTable(pool); const applied = await getAppliedMigrations(pool); const files = readdirSync(migrationsDir) .filter((f) => f.endsWith(".sql")) .sort(); // Alphabetical order: 001_create_users.sql, 002_add_email.sql const newMigrations: string[] = []; for (const file of files) { if (applied.has(file)) continue; const sql = readFileSync(join(migrationsDir, file), "utf-8"); const client = await pool.connect(); try { await client.query("BEGIN"); await client.query(sql); await client.query(`INSERT INTO ${MIGRATIONS_TABLE} (name) VALUES ($1)`, [ file, ]); await client.query("COMMIT"); newMigrations.push(file); } catch (err) { await client.query("ROLLBACK"); throw new Error( `Migration ${file} failed: ${err instanceof Error ? err.message : String(err)}`, ); } finally { client.release(); } } return newMigrations; } export { runMigrations }; ``` **Why good:** Each migration runs in its own transaction (atomic per file), tracks applied migrations in database, alphabetical ordering ensures deterministic execution, no external dependencies **When NOT to use:** If you already use a migration tool from your ORM or framework, use that instead. This pattern is for projects that use raw `pg` without higher-level tools. --- ## Testing with Transaction Rollback Wrap each test in a transaction that rolls back after the test. This provides test isolation without database cleanup scripts. ```typescript import type pg from "pg"; async function withTestTransaction<T>( pool: pg.Pool, testFn: (client: pg.PoolClient) => Promise<T>, ): Promise<T> { const client = await pool.connect(); try { await client.query("BEGIN"); const result = await testFn(client); return result; } finally { await client.query("ROLLBACK"); client.release(); } } // Usage in tests (framework-agnostic) // describe("UserRepository", () => { // it("creates a user", async () => { // await withTestTransaction(pool, async (client) => { // await client.query( // "INSERT INTO users (name, email) VALUES ($1, $2)", // ["Test User", "[email protected]"], // ); // // const { rows } = await client.query( // "SELECT * FROM users WHERE email = $1", // ["[email protected]"], // ); // // expect(rows).toHaveLength(1); // expect(rows[0].name).toBe("Test User"); // // ROLLBACK happens in finally -- database is clean for next test // }); // }); // }); ``` **Why good:** Every test starts with a clean database state, no teardown needed, tests are isolated, fast because rollback is cheaper than truncate **Gotcha:** The test function receives a `PoolClient`, not the pool itself. Code under test must accept a client/pool interface so you can inject the test client. If your code calls `pool.query()` internally, it will use a different connection and won't see the test transaction's uncommitted data. --- ## Connection Pool per Test Suite For integration tests that need `pool.query()` to work (not just a dedicated client), create an isolated pool per test suite with a test database. ```typescript import pg from "pg"; const POOL_MAX_TEST = 5; const IDLE_TIMEOUT_TEST_MS = 1_000; const CONNECTION_TIMEOUT_TEST_MS = 3_000; function createTestPool(): pg.Pool { const pool = new pg.Pool({ connectionString: process.env.TEST_DATABASE_URL, max: POOL_MAX_TEST, idleTimeoutMillis: IDLE_TIMEOUT_TEST_MS, connectionTimeoutMillis: CONNECTION_TIMEOUT_TEST_MS, allowExitOnIdle: true, // Let test runner exit cleanly }); pool.on("error", (err) => { console.error("Test pool error:", err.message); }); return pool; } // Truncate all tables between test suites async function cleanDatabase(pool: pg.Pool): Promise<void> { await pool.query(` DO $$ DECLARE r RECORD; BEGIN FOR r IN ( SELECT tablename FROM pg_tables WHERE schemaname = 'public' AND tablename != 'schema_migrations' ) LOOP EXECUTE 'TRUNCATE TABLE ' || quote_ident(r.tablename) || ' CASCADE'; END LOOP; END $$ `); } export { createTestPool, cleanDatabase }; ``` **Why good:** `allowExitOnIdle: true` prevents test runner from hanging, low pool size for tests, truncate cascades handle foreign keys, skips migration table --- _Full skill documentation: [SKILL.md](../SKILL.md) | Quick reference: [reference.md](../reference.md)_ -
core.md 7.3 KB
# PostgreSQL -- Core Pattern Examples > Pool setup, parameterized queries, typed results, and error handling. Reference from [SKILL.md](../SKILL.md). **Related examples:** - [transactions.md](transactions.md) -- BEGIN/COMMIT/ROLLBACK, savepoints, retry logic - [streaming.md](streaming.md) -- Cursors, pg-query-stream, LISTEN/NOTIFY - [advanced.md](advanced.md) -- SSL/TLS, prepared statements, testing patterns --- ## Pool Setup with Full Error Handling ```typescript import pg from "pg"; const POOL_MAX_CLIENTS = 20; const IDLE_TIMEOUT_MS = 30_000; const CONNECTION_TIMEOUT_MS = 5_000; const MAX_LIFETIME_SECONDS = 1_800; function createPool(): pg.Pool { const connectionString = process.env.DATABASE_URL; if (!connectionString) { throw new Error("DATABASE_URL environment variable is required"); } const pool = new pg.Pool({ connectionString, max: POOL_MAX_CLIENTS, idleTimeoutMillis: IDLE_TIMEOUT_MS, connectionTimeoutMillis: CONNECTION_TIMEOUT_MS, maxLifetimeSeconds: MAX_LIFETIME_SECONDS, }); // REQUIRED: idle client errors crash the process if unhandled pool.on("error", (err) => { console.error("Unexpected idle client error:", err.message); }); return pool; } export { createPool }; ``` **Why good:** Environment variable validation, named constants for all config values, error handler prevents process crash, `connectionTimeoutMillis` prevents infinite waits ```typescript // ❌ Bad Example - Hardcoded, no error handling import pg from "pg"; const pool = new pg.Pool({ host: "localhost", database: "mydb", user: "root", password: "secret123", }); // No pool.on("error") -- idle client errors crash process // Hardcoded credentials -- will be committed to version control // Default connectionTimeoutMillis: 0 -- infinite wait on pool exhaustion ``` **Why bad:** Hardcoded credentials leak in version control, missing error handler crashes process, default timeout means pool exhaustion hangs forever --- ## Parameterized Queries ```typescript import type pg from "pg"; interface UserRow { id: number; name: string; email: string; created_at: Date; } // Simple SELECT with parameters async function getUserById( pool: pg.Pool, userId: number, ): Promise<UserRow | undefined> { const result = await pool.query<UserRow>( "SELECT id, name, email, created_at FROM users WHERE id = $1", [userId], ); return result.rows[0]; } // INSERT with RETURNING async function createUser( pool: pg.Pool, name: string, email: string, ): Promise<UserRow> { const result = await pool.query<UserRow>( "INSERT INTO users (name, email) VALUES ($1, $2) RETURNING id, name, email, created_at", [name, email], ); return result.rows[0]; } // UPDATE with rowCount check async function updateUserEmail( pool: pg.Pool, userId: number, newEmail: string, ): Promise<boolean> { const result = await pool.query( "UPDATE users SET email = $1, updated_at = NOW() WHERE id = $2", [newEmail, userId], ); return (result.rowCount ?? 0) > 0; } // DELETE async function deleteUser(pool: pg.Pool, userId: number): Promise<boolean> { const result = await pool.query("DELETE FROM users WHERE id = $1", [userId]); return (result.rowCount ?? 0) > 0; } export { getUserById, createUser, updateUserEmail, deleteUser }; ``` **Why good:** Generic `<UserRow>` types result rows, `$1`/`$2` prevents injection, `RETURNING` avoids a second query, `rowCount` check for update/delete success --- ## Array Parameters PostgreSQL array parameters have a specific gotcha: you cannot use `IN ($1)` with a JavaScript array. Use `= ANY($1)` instead. ```typescript // ✅ Good Example - Array parameter with ANY async function getUsersByIds( pool: pg.Pool, userIds: number[], ): Promise<UserRow[]> { const result = await pool.query<UserRow>( "SELECT id, name, email FROM users WHERE id = ANY($1)", [userIds], // pg auto-converts JS array to PostgreSQL array ); return result.rows; } ``` **Why good:** `= ANY($1)` works with pg's automatic array conversion; no need to generate `$1, $2, $3, ...` placeholders dynamically ```typescript // ❌ Bad Example - IN with single parameter const result = await pool.query("SELECT * FROM users WHERE id IN ($1)", [ [1, 2, 3], ]); // Error: invalid input syntax for type integer: "{1,2,3}" // pg converts [1,2,3] to the string "{1,2,3}" which IN cannot parse ``` **Why bad:** `IN ($1)` expects a single scalar value, not an array; pg's array conversion creates a PostgreSQL array literal that IN cannot parse --- ## Error Handling with PostgreSQL Error Codes ```typescript import type pg from "pg"; // Named constants for PostgreSQL SQLSTATE error codes const PG_UNIQUE_VIOLATION = "23505"; const PG_FOREIGN_KEY_VIOLATION = "23503"; const PG_NOT_NULL_VIOLATION = "23502"; const PG_CHECK_VIOLATION = "23514"; const PG_DEADLOCK_DETECTED = "40P01"; const PG_SERIALIZATION_FAILURE = "40001"; // Type guard for PostgreSQL errors interface PgError extends Error { code: string; constraint?: string; detail?: string; table?: string; column?: string; schema?: string; } function isPgError(err: unknown): err is PgError { return err instanceof Error && "code" in err; } // Usage: map PostgreSQL errors to application errors async function createUser( pool: pg.Pool, email: string, name: string, ): Promise<UserRow> { try { const result = await pool.query<UserRow>( "INSERT INTO users (email, name) VALUES ($1, $2) RETURNING *", [email, name], ); return result.rows[0]; } catch (err) { if (!isPgError(err)) throw err; switch (err.code) { case PG_UNIQUE_VIOLATION: throw new ConflictError( `Duplicate value for constraint: ${err.constraint}`, ); case PG_FOREIGN_KEY_VIOLATION: throw new NotFoundError( `Referenced entity does not exist: ${err.detail}`, ); case PG_NOT_NULL_VIOLATION: throw new ValidationError(`Missing required field: ${err.column}`); case PG_CHECK_VIOLATION: throw new ValidationError(`Validation failed: ${err.constraint}`); default: throw err; } } } export { isPgError, PG_UNIQUE_VIOLATION, PG_FOREIGN_KEY_VIOLATION, PG_NOT_NULL_VIOLATION, PG_CHECK_VIOLATION, PG_DEADLOCK_DETECTED, PG_SERIALIZATION_FAILURE, }; ``` **Why good:** Named constants replace magic strings, type guard enables safe property access, switch on code maps to domain-specific errors, re-throws unknown errors --- ## Graceful Pool Shutdown ```typescript import type pg from "pg"; async function gracefulShutdown(pool: pg.Pool): Promise<void> { console.log("Shutting down database pool..."); await pool.end(); // Waits for checked-out clients to be released, then closes all console.log("Database pool closed"); } // Wire into process signals process.on("SIGTERM", async () => { await gracefulShutdown(pool); process.exit(0); }); process.on("SIGINT", async () => { await gracefulShutdown(pool); process.exit(0); }); ``` **Why good:** `pool.end()` drains the pool cleanly, signal handlers prevent abrupt disconnection **Gotcha:** `pool.end()` waits for ALL checked-out clients to be released. If any client was never released (leaked), `pool.end()` will hang forever. This is a symptom of missing `finally` blocks in `pool.connect()` usage. --- _Full skill documentation: [SKILL.md](../SKILL.md) | Quick reference: [reference.md](../reference.md)_ -
streaming.md 7.2 KB
# PostgreSQL -- Streaming & Real-Time Examples > Cursors, pg-query-stream for large result sets, and LISTEN/NOTIFY for real-time notifications. Reference from [SKILL.md](../SKILL.md). **Related examples:** - [core.md](core.md) -- Pool setup, parameterized queries, error handling - [transactions.md](transactions.md) -- Transactions, savepoints, retry logic - [advanced.md](advanced.md) -- Prepared statements, SSL, testing patterns --- ## Cursor-Based Batch Processing (pg-cursor) Use `pg-cursor` when you need to process rows in batches without loading the entire result set into memory. ```typescript import type pg from "pg"; import Cursor from "pg-cursor"; const BATCH_SIZE = 500; async function processLargeTable( pool: pg.Pool, status: string, processor: (rows: OrderRow[]) => Promise<void>, ): Promise<number> { const client = await pool.connect(); let totalProcessed = 0; try { const cursor = client.query( new Cursor( "SELECT id, user_id, total, status FROM orders WHERE status = $1 ORDER BY id", [status], ), ); let rows = await cursor.read(BATCH_SIZE); while (rows.length > 0) { await processor(rows); totalProcessed += rows.length; rows = await cursor.read(BATCH_SIZE); } await cursor.close(); return totalProcessed; } finally { client.release(); } } interface OrderRow { id: number; user_id: number; total: string; // numeric comes back as string status: string; } export { processLargeTable }; ``` **Why good:** Fixed memory footprint regardless of result set size, parameterized cursor query, proper client release, cursor closed explicitly **When to use:** Processing 1,000-100,000 rows where you need to control batch size and do async work per batch. --- ## Streaming with pg-query-stream Use `pg-query-stream` when you need a Node.js Readable stream -- ideal for piping to file writes, HTTP responses, or transform streams. ```typescript import type pg from "pg"; import QueryStream from "pg-query-stream"; import { pipeline } from "node:stream/promises"; import { createWriteStream } from "node:fs"; import { Transform } from "node:stream"; const STREAM_BATCH_SIZE = 100; async function exportOrdersToCsv( pool: pg.Pool, outputPath: string, ): Promise<void> { const client = await pool.connect(); try { const queryStream = new QueryStream( "SELECT id, user_id, total, created_at FROM orders ORDER BY id", [], { batchSize: STREAM_BATCH_SIZE }, ); const dbStream = client.query(queryStream); const csvTransform = new Transform({ objectMode: true, transform(row, _encoding, callback) { const line = `${row.id},${row.user_id},${row.total},${row.created_at.toISOString()}\n`; callback(null, line); }, }); const fileStream = createWriteStream(outputPath); // Write CSV header fileStream.write("id,user_id,total,created_at\n"); await pipeline(dbStream, csvTransform, fileStream); } finally { client.release(); } } export { exportOrdersToCsv }; ``` **Why good:** Constant memory usage for arbitrarily large tables, `pipeline` handles backpressure and error propagation, client released in finally **When to use:** Exporting 100,000+ rows to files, HTTP streaming responses, or ETL pipelines. For smaller result sets or when you need per-batch async processing, use pg-cursor instead. --- ## LISTEN/NOTIFY -- Subscribing to PostgreSQL Notifications LISTEN/NOTIFY is PostgreSQL's built-in pub/sub mechanism. A client executes `LISTEN channel` and receives notifications sent via `NOTIFY channel` or `pg_notify(channel, payload)`. ```typescript import type pg from "pg"; const ORDER_CHANNEL = "order_events"; interface OrderEvent { orderId: number; action: "created" | "updated" | "canceled"; timestamp: string; } async function subscribeToOrderEvents( pool: pg.Pool, handler: (event: OrderEvent) => void, ): Promise<{ unsubscribe: () => Promise<void> }> { // LISTEN requires a dedicated client -- it stays checked out const client = await pool.connect(); client.on("notification", (msg) => { if (msg.channel === ORDER_CHANNEL && msg.payload) { try { const event = JSON.parse(msg.payload) as OrderEvent; handler(event); } catch (err) { console.error("Failed to parse notification payload:", err); } } }); await client.query(`LISTEN ${ORDER_CHANNEL}`); return { unsubscribe: async () => { await client.query(`UNLISTEN ${ORDER_CHANNEL}`); client.release(); }, }; } export { subscribeToOrderEvents }; export type { OrderEvent }; ``` **Why good:** Dedicated client for listener, JSON payload parsing with error handling, unsubscribe releases client back to pool, typed event interface --- ## LISTEN/NOTIFY -- Publishing Notifications Use `pg_notify()` with parameterized queries for safe payload publishing. This can use `pool.query()` since it is a single command. ```typescript import type pg from "pg"; const ORDER_CHANNEL = "order_events"; async function publishOrderEvent( pool: pg.Pool, orderId: number, action: "created" | "updated" | "canceled", ): Promise<void> { const payload = JSON.stringify({ orderId, action, timestamp: new Date().toISOString(), }); // pg_notify() with parameterized channel and payload await pool.query("SELECT pg_notify($1, $2)", [ORDER_CHANNEL, payload]); } export { publishOrderEvent }; ``` **Why good:** `pg_notify($1, $2)` is parameterized (safe from injection), `pool.query()` is fine for publishing (no need for dedicated client), JSON payload is structured ```typescript // ❌ Bad Example - String interpolation in NOTIFY await pool.query(`NOTIFY ${channel}, '${payload}'`); // SQL injection if channel or payload contain quotes // Also: NOTIFY payload is limited to 8000 bytes and must be a string literal ``` **Why bad:** String interpolation allows injection, `NOTIFY` with literal payload has an 8000-byte limit and requires manual escaping --- ## LISTEN/NOTIFY -- Database Trigger Integration A common pattern is to fire NOTIFY from a PostgreSQL trigger so that application code is notified of data changes automatically. ```sql -- PostgreSQL trigger function CREATE OR REPLACE FUNCTION notify_order_change() RETURNS TRIGGER AS $$ BEGIN PERFORM pg_notify( 'order_events', json_build_object( 'orderId', NEW.id, 'action', TG_OP, 'timestamp', NOW() )::text ); RETURN NEW; END; $$ LANGUAGE plpgsql; -- Attach trigger to table CREATE TRIGGER order_change_trigger AFTER INSERT OR UPDATE ON orders FOR EACH ROW EXECUTE FUNCTION notify_order_change(); ``` **Why good:** Notifications are automatic -- no application code needed for publishing. Trigger fires within the transaction, so notification is only sent on COMMIT. **Gotcha:** If the transaction rolls back, the notification is NOT sent. This is usually desired behavior -- you don't want to notify about changes that didn't persist. **Gotcha:** LISTEN is NOT transactional -- `BEGIN; LISTEN channel; ROLLBACK;` still leaves the listener active. Only UNLISTEN or disconnecting removes it. --- _Full skill documentation: [SKILL.md](../SKILL.md) | Quick reference: [reference.md](../reference.md)_ -
transactions.md 8.8 KB
# PostgreSQL -- Transaction Examples > BEGIN/COMMIT/ROLLBACK, savepoints, retry logic for deadlocks, and advisory locks. Reference from [SKILL.md](../SKILL.md). **Related examples:** - [core.md](core.md) -- Pool setup, parameterized queries, error handling - [streaming.md](streaming.md) -- Cursors, pg-query-stream, LISTEN/NOTIFY - [advanced.md](advanced.md) -- Prepared statements, testing with transaction rollback --- ## Basic Transaction ```typescript import type pg from "pg"; interface TransferResult { fromBalance: number; toBalance: number; } async function transferFunds( pool: pg.Pool, fromAccountId: number, toAccountId: number, amount: number, ): Promise<TransferResult> { const client = await pool.connect(); try { await client.query("BEGIN"); // Lock rows in consistent order to prevent deadlocks const { rows } = await client.query<{ id: number; balance: number }>( `SELECT id, balance FROM accounts WHERE id = ANY($1) ORDER BY id FOR UPDATE`, [[fromAccountId, toAccountId]], ); const fromAccount = rows.find((r) => r.id === fromAccountId); const toAccount = rows.find((r) => r.id === toAccountId); if (!fromAccount || !toAccount) { throw new Error("Account not found"); } if (fromAccount.balance < amount) { throw new Error("Insufficient balance"); } await client.query( "UPDATE accounts SET balance = balance - $1 WHERE id = $2", [amount, fromAccountId], ); await client.query( "UPDATE accounts SET balance = balance + $1 WHERE id = $2", [amount, toAccountId], ); await client.query("COMMIT"); return { fromBalance: fromAccount.balance - amount, toBalance: toAccount.balance + amount, }; } catch (err) { await client.query("ROLLBACK"); throw err; } finally { client.release(); } } export { transferFunds }; ``` **Why good:** `FOR UPDATE` locks rows to prevent concurrent modification, consistent `ORDER BY id` prevents deadlocks, balance check after locking prevents race conditions, ROLLBACK on any error, release in finally --- ## Transaction Helper (Reusable) Extract the try/BEGIN/COMMIT/catch/ROLLBACK/finally/release boilerplate into a reusable helper. ```typescript import type pg from "pg"; async function withTransaction<T>( pool: pg.Pool, callback: (client: pg.PoolClient) => Promise<T>, ): Promise<T> { const client = await pool.connect(); try { await client.query("BEGIN"); const result = await callback(client); await client.query("COMMIT"); return result; } catch (err) { await client.query("ROLLBACK"); throw err; } finally { client.release(); } } export { withTransaction }; ``` Usage: ```typescript const order = await withTransaction(pool, async (client) => { const { rows: [newOrder], } = await client.query<OrderRow>( "INSERT INTO orders (user_id, total) VALUES ($1, $2) RETURNING *", [userId, total], ); for (const item of items) { await client.query( "INSERT INTO order_items (order_id, product_id, quantity) VALUES ($1, $2, $3)", [newOrder.id, item.productId, item.quantity], ); } return newOrder; }); ``` **Why good:** Eliminates boilerplate, impossible to forget ROLLBACK or release, type-safe return value --- ## Savepoints (Nested Transaction Behavior) PostgreSQL does not support nested transactions, but savepoints provide similar functionality within a transaction. ```typescript import type pg from "pg"; async function createOrderWithOptionalDiscount( pool: pg.Pool, userId: number, items: Array<{ productId: number; quantity: number }>, discountCode: string | null, ): Promise<void> { const client = await pool.connect(); try { await client.query("BEGIN"); // Core order creation const { rows: [order], } = await client.query<{ id: number }>( "INSERT INTO orders (user_id) VALUES ($1) RETURNING id", [userId], ); for (const item of items) { await client.query( "INSERT INTO order_items (order_id, product_id, quantity) VALUES ($1, $2, $3)", [order.id, item.productId, item.quantity], ); } // Optional: apply discount -- if it fails, order still succeeds if (discountCode) { await client.query("SAVEPOINT apply_discount"); try { await client.query( "UPDATE discount_codes SET used = true WHERE code = $1 AND used = false", [discountCode], ); await client.query( "UPDATE orders SET discount_code = $1 WHERE id = $2", [discountCode, order.id], ); } catch { // Discount failed -- rollback to savepoint, continue with order await client.query("ROLLBACK TO SAVEPOINT apply_discount"); } } await client.query("COMMIT"); } catch (err) { await client.query("ROLLBACK"); throw err; } finally { client.release(); } } export { createOrderWithOptionalDiscount }; ``` **Why good:** Savepoint isolates optional logic -- discount failure doesn't kill the order. `ROLLBACK TO SAVEPOINT` undoes only the savepoint's changes. The outer transaction can still COMMIT. --- ## Retry Logic for Deadlocks and Serialization Failures Deadlocks (`40P01`) and serialization failures (`40001`) are safe to retry. Wrap transactional code in a retry loop. ```typescript import type pg from "pg"; const PG_DEADLOCK_DETECTED = "40P01"; const PG_SERIALIZATION_FAILURE = "40001"; const MAX_RETRIES = 3; const BASE_DELAY_MS = 50; const RETRYABLE_CODES = new Set([ PG_DEADLOCK_DETECTED, PG_SERIALIZATION_FAILURE, ]); interface PgError extends Error { code: string; } function isPgError(err: unknown): err is PgError { return err instanceof Error && "code" in err; } async function withRetry<T>( pool: pg.Pool, operation: (client: pg.PoolClient) => Promise<T>, ): Promise<T> { for (let attempt = 0; attempt <= MAX_RETRIES; attempt++) { const client = await pool.connect(); try { await client.query("BEGIN"); const result = await operation(client); await client.query("COMMIT"); return result; } catch (err) { await client.query("ROLLBACK"); const isRetryable = isPgError(err) && RETRYABLE_CODES.has(err.code); if (!isRetryable || attempt === MAX_RETRIES) { throw err; } // Exponential backoff with jitter const delay = BASE_DELAY_MS * Math.pow(2, attempt) + Math.random() * BASE_DELAY_MS; await new Promise((resolve) => setTimeout(resolve, delay)); } finally { client.release(); } } // Unreachable, but TypeScript needs it throw new Error("Retry loop exited unexpectedly"); } export { withRetry }; ``` **Why good:** Only retries errors that are safe to retry (deadlocks, serialization failures), exponential backoff with jitter prevents thundering herd, bounded retries prevent infinite loops, fresh client per attempt --- ## Advisory Locks PostgreSQL advisory locks are application-level locks that don't lock any table rows. Use them to coordinate distributed operations. ```typescript import type pg from "pg"; // Advisory locks use a bigint key. Use a consistent hash for string-based lock names. const LOCK_REPORT_GENERATION = 1001; const LOCK_CACHE_REBUILD = 1002; async function withAdvisoryLock<T>( pool: pg.Pool, lockId: number, operation: (client: pg.PoolClient) => Promise<T>, ): Promise<T | null> { const client = await pool.connect(); try { // pg_try_advisory_lock returns true if lock acquired, false if already held const { rows } = await client.query<{ pg_try_advisory_lock: boolean }>( "SELECT pg_try_advisory_lock($1)", [lockId], ); if (!rows[0].pg_try_advisory_lock) { return null; // Another process holds the lock } try { return await operation(client); } finally { await client.query("SELECT pg_advisory_unlock($1)", [lockId]); } } finally { client.release(); } } // Usage const result = await withAdvisoryLock( pool, LOCK_REPORT_GENERATION, async (client) => { // Only one process runs this at a time across all app instances return generateExpensiveReport(client); }, ); if (result === null) { console.log("Report generation already in progress"); } export { withAdvisoryLock, LOCK_REPORT_GENERATION, LOCK_CACHE_REBUILD }; ``` **Why good:** `pg_try_advisory_lock` is non-blocking (returns false immediately if held), lock is released in a nested finally, named constants for lock IDs prevent collisions **Gotcha:** Advisory locks are bound to the **session** (connection), not the transaction. If the client disconnects, the lock is automatically released. If you use `pg_advisory_xact_lock`, the lock is released when the transaction ends instead. --- _Full skill documentation: [SKILL.md](../SKILL.md) | Quick reference: [reference.md](../reference.md)_
-
-
reference.md 12.4 KB
# PostgreSQL (pg) Quick Reference > Pool configuration, QueryResult properties, PostgreSQL error codes, type mapping, and production checklist. See [SKILL.md](SKILL.md) for core concepts and [examples/](examples/) for code examples. --- ## Pool Configuration Options | Option | Type | Default | Description | | ------------------------- | -------------------------- | ----------- | ------------------------------------------------------------------- | | `connectionString` | string | — | PostgreSQL connection URI (`postgresql://user:pass@host:5432/db`) | | `host` | string | `localhost` | Server hostname (or env `PGHOST`) | | `port` | number | `5432` | Server port (or env `PGPORT`) | | `database` | string | — | Database name (or env `PGDATABASE`) | | `user` | string | — | Username (or env `PGUSER`) | | `password` | string \| function | — | Password or async callback (or env `PGPASSWORD`) | | `max` | number | `10` | Maximum clients in the pool | | `min` | number | `0` | Minimum idle clients to retain | | `idleTimeoutMillis` | number | `10000` | Idle time before client is disconnected (0 = disabled) | | `connectionTimeoutMillis` | number | `0` | Timeout for new connections (0 = no timeout) | | `maxUses` | number | `Infinity` | Max times a client can be checked out before replacement | | `maxLifetimeSeconds` | number | `0` | Max age for connections (0 = disabled) | | `allowExitOnIdle` | boolean | `false` | Let event loop exit when all clients idle | | `ssl` | boolean \| TlsOptions | — | SSL/TLS configuration (see SSL section) | | `onConnect` | function \| async function | — | Setup callback invoked once per new client before it joins the pool | ### Recommended Production Configuration ```typescript const POOL_MAX = 20; const IDLE_TIMEOUT_MS = 30_000; const CONNECTION_TIMEOUT_MS = 5_000; const MAX_LIFETIME_SECONDS = 1800; // 30 minutes const pool = new pg.Pool({ connectionString: process.env.DATABASE_URL, max: POOL_MAX, idleTimeoutMillis: IDLE_TIMEOUT_MS, connectionTimeoutMillis: CONNECTION_TIMEOUT_MS, maxLifetimeSeconds: MAX_LIFETIME_SECONDS, ssl: process.env.NODE_ENV === "production" ? { rejectUnauthorized: true } : undefined, }); ``` ### Recommended Test Configuration ```typescript const pool = new pg.Pool({ connectionString: process.env.TEST_DATABASE_URL, max: 5, idleTimeoutMillis: 1_000, connectionTimeoutMillis: 3_000, allowExitOnIdle: true, }); ``` --- ## Pool Properties | Property | Type | Description | | ------------------- | ------ | ---------------------------------------------- | | `pool.totalCount` | number | Total clients in the pool (idle + checked out) | | `pool.idleCount` | number | Clients currently idle | | `pool.waitingCount` | number | Queued requests waiting for a client | --- ## Pool Events | Event | Callback Signature | Description | | --------- | ------------------------------ | --------------------------------------- | | `error` | `(err: Error, client: Client)` | Idle client backend error (MUST handle) | | `connect` | `(client: Client)` | New client connected to backend | | `acquire` | `(client: Client)` | Client checked out from pool | | `release` | `(err: Error, client: Client)` | Client returned to pool | | `remove` | `(client: Client)` | Client disconnected and removed | --- ## QueryResult Properties | Property | Type | Description | | ----------------- | ---------------- | ----------------------------------------------------------------------------- | | `result.rows` | `T[]` | Array of row objects (typed via generic) | | `result.rowCount` | `number \| null` | Rows affected by INSERT/UPDATE/DELETE. For SELECT, use `rows.length` instead. | | `result.command` | `string` | SQL command executed (`SELECT`, `INSERT`, `UPDATE`, `DELETE`, etc.) | | `result.fields` | `FieldInfo[]` | Column metadata (name, dataTypeID) | --- ## PostgreSQL Error Codes (Common) ### Constraint Violations (Class 23) | Code | Name | Typical HTTP | When It Fires | | ------- | ----------------------- | --------------- | -------------------------------------------- | | `23505` | `unique_violation` | 409 Conflict | INSERT/UPDATE violates UNIQUE or PRIMARY KEY | | `23503` | `foreign_key_violation` | 400 Bad Request | Referenced row does not exist | | `23502` | `not_null_violation` | 400 Bad Request | NULL in a NOT NULL column | | `23514` | `check_violation` | 400 Bad Request | CHECK constraint failed | ### Transaction / Concurrency (Class 40) | Code | Name | Action | When It Fires | | ------- | ----------------------- | ------------------ | --------------------------------------------- | | `40P01` | `deadlock_detected` | Retry with backoff | Two transactions waiting on each other | | `40001` | `serialization_failure` | Retry with backoff | Concurrent SERIALIZABLE transactions conflict | ### Connection (Class 08) | Code | Name | Action | When It Fires | | ------- | --------------------------- | ------------ | --------------------------- | | `08000` | `connection_exception` | Pool handles | General connection failure | | `08003` | `connection_does_not_exist` | Pool handles | Client disconnected | | `08006` | `connection_failure` | Pool handles | Could not connect to server | ### Other | Code | Name | Action | When It Fires | | ------- | ----------------------------- | ------------------------- | ----------------------------------------------------- | | `57014` | `query_canceled` | Check `statement_timeout` | Query exceeded timeout | | `42P01` | `undefined_table` | Fix SQL | Table does not exist | | `42703` | `undefined_column` | Fix SQL | Column does not exist | | `22P02` | `invalid_text_representation` | Fix input | Invalid input for data type (e.g., string to integer) | --- ## PostgreSQL-to-JavaScript Type Mapping | PostgreSQL Type | JavaScript Type | Notes | | ------------------------------- | ---------------------- | ------------------------------------------------------------ | | `integer`, `smallint`, `bigint` | `number` / `string` | `bigint` returns as string (exceeds JS Number range) | | `real`, `double precision` | `number` | Standard JS float | | `numeric`, `decimal` | **`string`** | Preserved as string to avoid precision loss | | `boolean` | `boolean` | | | `text`, `varchar`, `char` | `string` | | | `timestamp`, `timestamptz` | `Date` | Parsed to JS Date object | | `date` | `Date` | Parsed to JS Date object | | `json`, `jsonb` | `object` | Auto-parsed from JSON | | `uuid` | `string` | | | `bytea` | `Buffer` | Binary data | | `interval` | `object` | Parsed to `{ years, months, days, hours, minutes, seconds }` | | `integer[]`, `text[]` | `number[]`, `string[]` | Parsed to JS arrays | **Gotcha:** `numeric`/`decimal` returns as **string** by default. This is intentional -- JavaScript `number` cannot represent arbitrary-precision decimals. Parse manually if you need a number: `parseFloat(row.price)`. **Gotcha:** `bigint` returns as **string** when the value exceeds `Number.MAX_SAFE_INTEGER`. Use `BigInt(row.id)` for arithmetic. --- ## SSL/TLS Quick Reference See [examples/advanced.md](examples/advanced.md) for full SSL/TLS code examples (cloud, mTLS, dynamic passwords). | Scenario | `ssl` value | Notes | | ------------------------------ | --------------------------------------------- | ------------------------------- | | Cloud-managed (RDS, Cloud SQL) | `{ rejectUnauthorized: true }` | Trusted CA, just enable | | Self-signed / mTLS | `{ ca, key, cert, rejectUnauthorized: true }` | Provide cert files via env vars | | Development only | `{ rejectUnauthorized: false }` | Never in production | **Gotcha:** If the connection string contains any SSL parameters (`sslmode`, `sslcert`, `sslkey`, `sslrootcert`), the entire `ssl` config object is replaced and any additional options are lost. Use one or the other, not both. **Gotcha:** `ssl: true` relies on Node.js TLS defaults (which reject unauthorized certs), but does not explicitly set `rejectUnauthorized`. Always use the object form (`ssl: { rejectUnauthorized: true }`) to be explicit about intent. --- ## Production Checklist ### Connection Management - [ ] Pool `error` event handler on every pool instance - [ ] `connectionTimeoutMillis` set (not default 0 = infinite wait) - [ ] `maxLifetimeSeconds` set to rotate connections (helps with DNS changes, PgBouncer) - [ ] `pool.end()` called on graceful shutdown (SIGTERM handler) - [ ] All `pool.connect()` calls release clients in `finally` blocks ### Security - [ ] Parameterized queries for ALL user input (`$1`, `$2`, ...) - [ ] SSL enabled in production (`ssl: { rejectUnauthorized: true }`) - [ ] Connection credentials from environment variables (not hardcoded) - [ ] Least-privilege database user (not superuser) - [ ] `statement_timeout` configured to prevent runaway queries ### Data Integrity - [ ] Transactions for multi-statement operations - [ ] Constraint violation errors handled (23505, 23503, etc.) - [ ] Deadlock/serialization errors handled with retry logic (40P01, 40001) - [ ] `numeric`/`decimal` values handled as strings (not silently converted to floats) ### Performance - [ ] Pool `max` sized appropriately (rule of thumb: 2-4x CPU cores) - [ ] Streaming for result sets > 1,000 rows - [ ] No `SELECT *` in production queries - [ ] Indexes on frequently queried columns - [ ] `EXPLAIN ANALYZE` for slow queries ### Monitoring - [ ] Track `pool.totalCount`, `pool.idleCount`, `pool.waitingCount` - [ ] Alert when `waitingCount > 0` sustained (pool exhaustion) - [ ] Log slow queries (`log_min_duration_statement` in postgresql.conf) - [ ] Monitor connection count against PostgreSQL `max_connections` --- _Full skill documentation: [SKILL.md](SKILL.md) | Examples: [examples/](examples/)_ -
SKILL.md 17.6 KB
--- name: api-database-postgresql description: Direct PostgreSQL access with node-postgres (pg) -- connection pools, parameterized queries, transactions, streaming, LISTEN/NOTIFY, error handling --- # PostgreSQL Patterns (node-postgres) > **Quick Guide:** Use the `pg` package (v8.x) for direct PostgreSQL access. **Always use `Pool`** -- never create individual `Client` instances in application code. Use **parameterized queries** (`$1`, `$2`) for ALL user input -- never interpolate strings into SQL. For transactions, check out a dedicated client with `pool.connect()` and use `BEGIN`/`COMMIT`/`ROLLBACK` in a `try`/`catch`/`finally` that always calls `client.release()`. Handle the pool `error` event to prevent process crashes from idle client errors. Use `pg-query-stream` for large result sets to avoid loading everything into memory. --- <critical_requirements> ## CRITICAL: Before Using This Skill > **All code must follow project conventions in CLAUDE.md** (kebab-case, named exports, import ordering, `import type`, named constants) **(You MUST use parameterized queries (`$1`, `$2`, ...) for ALL values -- NEVER concatenate or interpolate user input into SQL strings)** **(You MUST use `Pool` for all database access -- NEVER create standalone `Client` instances in application code)** **(You MUST release clients back to the pool in a `finally` block after `pool.connect()` -- leaked clients exhaust the pool and hang the application)** **(You MUST handle the `error` event on Pool instances -- unhandled idle client errors crash the Node.js process)** </critical_requirements> --- ## Examples - [Core Patterns](examples/core.md) -- Pool setup, parameterized queries, type-safe results, error handling - [Transactions](examples/transactions.md) -- BEGIN/COMMIT/ROLLBACK, savepoints, retry logic, advisory locks - [Streaming](examples/streaming.md) -- Cursors, pg-query-stream, LISTEN/NOTIFY for real-time - [Advanced](examples/advanced.md) -- SSL/TLS, prepared statements, migrations, testing patterns **Additional resources:** - [reference.md](reference.md) -- Pool options, error codes, QueryResult properties, production checklist --- **Auto-detection:** PostgreSQL, pg, node-postgres, Pool, Client, pool.query, pool.connect, client.query, $1, parameterized query, BEGIN, COMMIT, ROLLBACK, LISTEN, NOTIFY, pg_notify, pg-query-stream, pg-cursor, Cursor, QueryResult, QueryResultRow, connectionString, PGHOST, PGDATABASE, unique_violation, 23505, deadlock, 40P01, advisory lock **When to use:** - Direct SQL queries against PostgreSQL (not behind an ORM) - Connection pool management for Node.js/PostgreSQL applications - Transactions spanning multiple queries that must be atomic - Streaming large result sets without loading everything into memory - Real-time change notifications via LISTEN/NOTIFY - Integration testing with transaction rollback isolation **Key patterns covered:** - Pool configuration and lifecycle (creation, error handling, graceful shutdown) - Parameterized queries with `$1`-style placeholders (SQL injection prevention) - Type-safe query results using TypeScript generics - Transaction management with dedicated clients - Streaming with pg-cursor and pg-query-stream - LISTEN/NOTIFY for real-time PostgreSQL event handling - PostgreSQL error code handling (constraint violations, deadlocks, serialization failures) - SSL/TLS connection configuration - Testing with transaction rollback isolation **When NOT to use:** - You need an ORM or query builder -- use your ORM/query builder skill instead - You need in-memory caching -- use a caching solution - You need document storage without relational constraints -- use a document database - Simple key-value lookups at sub-millisecond latency -- use an in-memory data store --- <philosophy> ## Philosophy `pg` (node-postgres) is a **low-level PostgreSQL client** that gives you full control over SQL, connections, and transactions. The core principle: **write SQL directly, let PostgreSQL do the heavy lifting.** **Core principles:** 1. **Pool, never Client** -- Application code should always use `Pool`. The pool manages connections, handles reconnection, and prevents connection exhaustion. Use `pool.query()` for single queries, `pool.connect()` when you need a dedicated client (transactions). 2. **Parameterized everything** -- Never build SQL by string concatenation. Use `$1`, `$2` placeholders. This prevents SQL injection AND enables PostgreSQL query plan caching. 3. **Release in finally** -- Any client obtained via `pool.connect()` must be released in a `finally` block. A leaked client sits checked out forever, and once `max` clients leak, the pool deadlocks. 4. **Fail loudly** -- Handle the pool's `error` event. Handle query errors with specific PostgreSQL error codes. Never swallow errors silently. 5. **Stream large results** -- Don't `SELECT *` a million rows into memory. Use `pg-cursor` or `pg-query-stream` for large result sets. </philosophy> --- <patterns> ## Core Patterns ### Pattern 1: Pool Setup Create a single pool per database at application startup. See [examples/core.md](examples/core.md) for full configuration examples. ```typescript // ✅ Good Example - Pool with error handling import pg from "pg"; const POOL_MAX_CLIENTS = 20; const IDLE_TIMEOUT_MS = 30_000; const CONNECTION_TIMEOUT_MS = 5_000; function createPool(): pg.Pool { const pool = new pg.Pool({ connectionString: process.env.DATABASE_URL, max: POOL_MAX_CLIENTS, idleTimeoutMillis: IDLE_TIMEOUT_MS, connectionTimeoutMillis: CONNECTION_TIMEOUT_MS, }); pool.on("error", (err) => { console.error("Unexpected idle client error:", err.message); }); return pool; } export { createPool }; ``` **Why good:** Named constants for pool config, environment variable for connection string, error handler prevents process crash from idle client errors ```typescript // ❌ Bad Example - No pool, standalone client import pg from "pg"; const client = new pg.Client("postgres://localhost/mydb"); await client.connect(); // One connection for entire app -- no pooling, no reconnection, // no concurrency. If client disconnects, app crashes. ``` **Why bad:** Standalone Client has no connection pooling, no automatic reconnection, no concurrency -- every query blocks on a single connection --- ### Pattern 2: Parameterized Queries Always use `$1`-style placeholders. See [examples/core.md](examples/core.md) for typed query helpers. ```typescript // ✅ Good Example - Parameterized query with typed result interface UserRow { id: number; name: string; email: string; } const result = await pool.query<UserRow>( "SELECT id, name, email FROM users WHERE id = $1", [userId], ); const user = result.rows[0]; // UserRow | undefined ``` **Why good:** `$1` placeholder prevents SQL injection, generic `<UserRow>` types the `rows` array, result is properly typed ```typescript // ❌ Bad Example - String interpolation (SQL INJECTION!) const result = await pool.query( `SELECT * FROM users WHERE name = '${userName}'`, ); // userName = "'; DROP TABLE users; --" -> catastrophic ``` **Why bad:** String interpolation allows SQL injection, no type safety on result rows, `SELECT *` returns untyped columns --- ### Pattern 3: Transactions Use `pool.connect()` to get a dedicated client for the transaction. See [examples/transactions.md](examples/transactions.md) for savepoints, retries, and advisory locks. ```typescript // ✅ Good Example - Transaction with proper cleanup async function transferFunds( pool: pg.Pool, fromId: number, toId: number, amount: number, ): Promise<void> { const client = await pool.connect(); try { await client.query("BEGIN"); await client.query( "UPDATE accounts SET balance = balance - $1 WHERE id = $2", [amount, fromId], ); await client.query( "UPDATE accounts SET balance = balance + $1 WHERE id = $2", [amount, toId], ); await client.query("COMMIT"); } catch (err) { await client.query("ROLLBACK"); throw err; } finally { client.release(); } } ``` **Why good:** Dedicated client via `pool.connect()`, `ROLLBACK` on error, `client.release()` in `finally` guarantees the client returns to the pool ```typescript // ❌ Bad Example - Transaction with pool.query() await pool.query("BEGIN"); await pool.query("UPDATE accounts SET balance = balance - 100 WHERE id = 1"); await pool.query("UPDATE accounts SET balance = balance + 100 WHERE id = 2"); await pool.query("COMMIT"); // Each pool.query() may use a DIFFERENT client -- the BEGIN/COMMIT // execute on different connections, so there is no transaction at all ``` **Why bad:** `pool.query()` checks out a random client each time -- BEGIN, UPDATEs, and COMMIT may run on different connections, so there is no actual transaction --- ### Pattern 4: Error Handling with PostgreSQL Error Codes PostgreSQL errors include a `code` field with the SQLSTATE error code. See [reference.md](reference.md) for the full error code table. ```typescript // ✅ Good Example - Handling specific PostgreSQL errors const PG_UNIQUE_VIOLATION = "23505"; const PG_FOREIGN_KEY_VIOLATION = "23503"; const PG_DEADLOCK_DETECTED = "40P01"; const PG_SERIALIZATION_FAILURE = "40001"; interface PgError extends Error { code: string; constraint?: string; detail?: string; table?: string; column?: string; } function isPgError(err: unknown): err is PgError { return err instanceof Error && "code" in err; } try { await pool.query("INSERT INTO users (email) VALUES ($1)", [email]); } catch (err) { if (isPgError(err) && err.code === PG_UNIQUE_VIOLATION) { throw new ConflictError(`Email already exists: ${err.constraint}`); } if (isPgError(err) && err.code === PG_DEADLOCK_DETECTED) { // Retry the operation } throw err; } ``` **Why good:** Named constants for error codes (no magic strings), type guard for safe property access, specific handling per error type, re-throws unknown errors --- ### Pattern 5: Streaming Large Result Sets Use `pg-cursor` or `pg-query-stream` for queries returning many rows. See [examples/streaming.md](examples/streaming.md) for full streaming patterns. ```typescript // ✅ Good Example - Cursor for batch processing import Cursor from "pg-cursor"; const BATCH_SIZE = 100; async function processAllOrders(pool: pg.Pool): Promise<void> { const client = await pool.connect(); try { const cursor = client.query( new Cursor("SELECT * FROM orders WHERE status = $1", ["pending"]), ); let rows = await cursor.read(BATCH_SIZE); while (rows.length > 0) { await processBatch(rows); rows = await cursor.read(BATCH_SIZE); } await cursor.close(); } finally { client.release(); } } ``` **Why good:** Processes rows in fixed-size batches without loading entire result set into memory, proper client release in `finally` --- ### Pattern 6: LISTEN/NOTIFY PostgreSQL can push real-time notifications to connected clients. See [examples/streaming.md](examples/streaming.md) for full examples. ```typescript // ✅ Good Example - LISTEN/NOTIFY with dedicated client const CHANNEL = "order_updates"; async function listenForUpdates(pool: pg.Pool): Promise<pg.PoolClient> { const client = await pool.connect(); client.on("notification", (msg) => { if (msg.channel === CHANNEL && msg.payload) { const data = JSON.parse(msg.payload); handleOrderUpdate(data); } }); await client.query(`LISTEN ${CHANNEL}`); return client; // Caller is responsible for release on shutdown } // Publishing from another connection await pool.query("SELECT pg_notify($1, $2)", [CHANNEL, JSON.stringify(data)]); ``` **Why good:** Dedicated client stays checked out for the lifetime of the listener, `pg_notify()` with parameterized channel/payload prevents injection, JSON payload for structured data **When to use:** Real-time notifications where sub-second latency matters and the volume is low-to-moderate (hundreds per second). For high-throughput streaming, use a dedicated message broker. </patterns> --- <decision_framework> ## Decision Framework ### pool.query() vs pool.connect() ``` Do I need a dedicated client? ├─ Single query, no transaction? -> pool.query() (auto-releases) ├─ Multiple queries in a transaction? -> pool.connect() + BEGIN/COMMIT/ROLLBACK ├─ LISTEN for notifications? -> pool.connect() (keep client for lifetime of listener) ├─ Cursor/streaming? -> pool.connect() (cursor binds to a connection) └─ Prepared statements across queries? -> pool.connect() (plan caches per connection) ``` ### Error Handling Strategy ``` What kind of PostgreSQL error? ├─ 23505 (unique_violation)? -> Map to 409 Conflict, include constraint name ├─ 23503 (foreign_key_violation)? -> Map to 400 Bad Request, entity not found ├─ 23502 (not_null_violation)? -> Map to 400 Bad Request, missing required field ├─ 23514 (check_violation)? -> Map to 400 Bad Request, validation failed ├─ 40P01 (deadlock_detected)? -> Retry with backoff (safe to retry) ├─ 40001 (serialization_failure)? -> Retry with backoff (safe to retry) ├─ 57014 (query_canceled)? -> Timeout, consider increasing statement_timeout ├─ 08xxx (connection_exception)? -> Pool handles reconnection, log and retry └─ Other? -> Log full error, return 500 ``` ### Streaming Decision ``` How many rows will the query return? ├─ < 1,000 rows? -> pool.query() is fine (result fits in memory) ├─ 1,000 - 100,000 rows? -> pg-cursor with batch processing ├─ 100,000+ rows? -> pg-query-stream piped to a writable stream └─ Need to export to file? -> pg-query-stream piped to file write stream ``` </decision_framework> --- <red_flags> ## RED FLAGS **High Priority Issues:** - Using string interpolation/concatenation for SQL values -- this is SQL injection, the most dangerous vulnerability in database code - Using `pool.query()` for transactions -- each call may use a different connection, so BEGIN/COMMIT have no effect - Not releasing clients after `pool.connect()` -- leaked clients exhaust the pool; once `max` clients leak, the app deadlocks on `pool.connect()` - Missing `pool.on("error")` handler -- idle client errors are emitted on the pool; unhandled, they crash the Node.js process - Using standalone `Client` in application code -- no pooling, no reconnection, no concurrency **Medium Priority Issues:** - `SELECT *` in production queries -- returns unnecessary columns, breaks when schema changes, prevents index-only scans - Loading millions of rows with `pool.query()` instead of streaming -- causes memory exhaustion and GC pressure - Hardcoded connection strings -- prevents environment-specific configuration, risks credential leaks in version control - Not handling specific PostgreSQL error codes -- generic error handling loses valuable information (which constraint, which column) - Using `LISTEN` with `pool.query()` -- notifications bind to a specific connection; pool.query releases the connection immediately **Common Mistakes:** - Forgetting that `result.rows[0]` can be `undefined` when no rows match -- always check before accessing - Relying on `result.rowCount` for `SELECT` emptiness checks -- use `result.rows.length` instead; `rowCount` is `null` for some commands (e.g., `LOCK`) and `rows.length` is universally reliable - Using `$1` inside string literals in SQL -- `'$1'` is a literal string, not a parameter; use `$1` outside quotes - Forgetting that PostgreSQL arrays in parameters are automatically converted -- `[1, 2, 3]` becomes `{1,2,3}` which works for `= ANY($1)` but not for `IN ($1)` (use `= ANY($1::int[])` instead of `IN`) - Calling `client.release(true)` routinely -- passing `true` destroys the client instead of returning it to the pool; only use after unrecoverable errors **Gotchas & Edge Cases:** - Pool `error` event vs query errors: Pool `error` fires for idle client backend disconnections (e.g., server restart). Query errors are thrown/rejected from the query call itself. You need both handlers. - `connectionTimeoutMillis: 0` (default) means **no timeout** -- connections wait forever if the pool is exhausted. Always set a timeout in production. - `idleTimeoutMillis` only affects clients that have been returned to the pool -- a checked-out client that is never released will never be cleaned up. - PostgreSQL `numeric`/`decimal` types are returned as **strings** by default (to avoid JavaScript floating-point precision loss). Parse them explicitly if you need numbers. - `LISTEN` survives transactions -- if you `BEGIN`, `LISTEN channel`, `ROLLBACK`, the listener is still active. LISTEN is not transactional. - `pool.end()` waits for all checked-out clients to be released. If a client is leaked (never released), `pool.end()` hangs forever. - SSL connections: if the connection string contains any SSL parameters (`sslmode`, `sslcert`, `sslkey`, `sslrootcert`), the entire `ssl` config object is replaced -- use one or the other, not both. </red_flags> --- <critical_reminders> ## CRITICAL REMINDERS > **All code must follow project conventions in CLAUDE.md** (kebab-case, named exports, import ordering, `import type`, named constants) **(You MUST use parameterized queries (`$1`, `$2`, ...) for ALL values -- NEVER concatenate or interpolate user input into SQL strings)** **(You MUST use `Pool` for all database access -- NEVER create standalone `Client` instances in application code)** **(You MUST release clients back to the pool in a `finally` block after `pool.connect()` -- leaked clients exhaust the pool and hang the application)** **(You MUST handle the `error` event on Pool instances -- unhandled idle client errors crash the Node.js process)** **Failure to follow these rules will cause SQL injection vulnerabilities, connection pool exhaustion, application hangs, and process crashes.** </critical_reminders>
Comments (0)
Sign in to join the conversation.
Reviews (0)
No reviews yet.
No comments yet.