/** * Capnweb-enabled PostgreSQL client for postgres.do * * This client provides: * - Magic map support for N+1 query elimination * - Promise pipelining for batched operations * - RPC transport with automatic batching * * @example * ```typescript * import { createClient } from 'postgres.do/rpc' * * const db = createClient('postgres.do/my-database') * * // Magic map - N+1 solved in single round trip * const usersWithOrders = await db.query('SELECT * FROM users') * .map(users => Promise.all( * users.map(u => db.query('SELECT * FROM orders WHERE user_id = $1', [u.id])) * )) * * // Promise pipelining - chain without await * const user = db.query('SELECT * FROM users WHERE id = $1', [id]) * const posts = user.map(u => db.query('SELECT * FROM posts WHERE author = $1', [u.id])) * const comments = posts.map(p => db.query('SELECT * FROM comments WHERE post_id = ANY($1)', [p.map(x => x.id)])) * const result = await comments // Single round trip! * * // Transactions with pipelining * await db.transaction(tx => { * const balance = tx.query('SELECT balance FROM accounts WHERE id = $1', [fromId]) * const updated = balance.map(b => { * if (b.balance < amount) throw new Error('Insufficient funds') * return tx.query('UPDATE accounts SET balance = balance - $1 WHERE id = $2', [amount, fromId]) * }) * return updated.map(() => * tx.query('UPDATE accounts SET balance = balance + $1 WHERE id = $2', [amount, toId]) * ) * }) * ``` */ import type { Row, IsolationLevel } from '../types' import { RpcTransport, createRpcTransport, type RpcTransportConfig } from '../transport/rpc' import { RpcPromise, createRpcPromise, magicMap, batchExecute, type RpcQueryResult, } from './rpc-promise' /** * Client configuration options */ export interface RpcClientConfig { /** Database URL (postgres.do/database-name or wss://...) */ url?: string /** API key for authentication */ apiKey?: string /** Custom WebSocket implementation */ WebSocket?: typeof WebSocket /** Connection timeout in milliseconds */ connectTimeout?: number /** Request timeout in milliseconds */ requestTimeout?: number /** Enable automatic batching (default: true) */ autoBatch?: boolean /** Custom transport configuration */ transport?: Partial } /** * Transaction options */ export interface TransactionOptions { /** Isolation level for the transaction */ isolationLevel?: IsolationLevel /** Read-only transaction hint */ readOnly?: boolean } /** * Transaction client interface */ export interface TransactionClient { /** Execute a query within the transaction */ query(sql: string, params?: unknown[]): RpcPromise /** Execute a query and return only the first row */ queryOne(sql: string, params?: unknown[]): RpcPromise /** Execute a query and return a scalar value */ queryScalar(sql: string, params?: unknown[]): RpcPromise /** Execute a statement that doesn't return rows */ execute(sql: string, params?: unknown[]): RpcPromise<{ rowCount: number }> } /** * Capnweb-enabled PostgreSQL client */ export class RpcClient { private transport: RpcTransport constructor(config: RpcClientConfig = {}) { const baseUrl = this.resolveUrl(config.url) const transportConfig: RpcTransportConfig = { url: baseUrl, autoBatch: config.autoBatch ?? true, } if (config.apiKey) transportConfig.apiKey = config.apiKey if (config.WebSocket) transportConfig.WebSocket = config.WebSocket if (config.connectTimeout) transportConfig.connectTimeout = config.connectTimeout if (config.requestTimeout) transportConfig.requestTimeout = config.requestTimeout if (config.transport?.batchDelayMs !== undefined) transportConfig.batchDelayMs = config.transport.batchDelayMs this.transport = createRpcTransport(transportConfig) } /** * Resolve the database URL */ private resolveUrl(url?: string): string { if (!url) { return 'wss://db.postgres.do/rpc' } // Handle postgres.do/database-name format if (url.startsWith('postgres.do/')) { const dbName = url.slice('postgres.do/'.length) return `wss://db.postgres.do/${dbName}/rpc` } // Handle full URLs if (url.startsWith('wss://') || url.startsWith('ws://')) { return url } // Handle postgres:// URLs if (url.startsWith('postgres://') || url.startsWith('postgresql://')) { const parsed = new URL(url) const protocol = parsed.protocol === 'postgres:' || parsed.protocol === 'postgresql:' ? 'wss' : 'ws' return `${protocol}://${parsed.host}${parsed.pathname}/rpc` } // Default: treat as database name return `wss://db.postgres.do/${url}/rpc` } /** * Execute a SQL query * * Returns an RpcPromise that supports .map() for magic map functionality. * * @example * ```typescript * // Simple query * const users = await db.query('SELECT * FROM users') * * // With magic map (N+1 elimination) * const usersWithOrders = await db.query('SELECT id FROM users') * .map(users => Promise.all(users.map(u => * db.query('SELECT * FROM orders WHERE user_id = $1', [u.id]) * ))) * ``` */ query(sql: string, params?: unknown[]): RpcPromise { return this.transport.queryRpc(sql, params) } /** * Execute a query and return only the first row * * @example * ```typescript * const user = await db.queryOne('SELECT * FROM users WHERE id = $1', [id]) * if (user) { * console.log(user.name) * } * ``` */ queryOne(sql: string, params?: unknown[]): RpcPromise { return createRpcPromise(async () => { const rows = await this.query(sql, params) return rows[0] ?? null }) } /** * Execute a query and return a scalar value (first column of first row) * * @example * ```typescript * const count = await db.queryScalar('SELECT COUNT(*) FROM users') * ``` */ queryScalar(sql: string, params?: unknown[]): RpcPromise { return createRpcPromise(async () => { const rows = await this.query(sql, params) if (!rows[0]) return null const firstKey = Object.keys(rows[0])[0] return (firstKey ? rows[0][firstKey] : null) as T }) } /** * Execute a SQL statement that doesn't return rows * * @example * ```typescript * const { rowCount } = await db.execute('DELETE FROM users WHERE id = $1', [id]) * console.log(`Deleted ${rowCount} rows`) * ``` */ execute(sql: string, params?: unknown[]): RpcPromise<{ rowCount: number }> { return createRpcPromise(async () => { const result = await this.transport.query(sql, params) return { rowCount: result.rowCount } }) } /** * Execute a batch of queries in a single round trip * * @example * ```typescript * const [users, orders, products] = await db.batch([ * { sql: 'SELECT * FROM users' }, * { sql: 'SELECT * FROM orders' }, * { sql: 'SELECT * FROM products' }, * ]) * ``` */ async batch( queries: Array<{ sql: string; params?: unknown[] }> ): Promise { const results = await this.transport.batch(queries) return results.map(r => r.rows) } /** * Execute queries within a transaction * * The transaction callback receives a transaction client that supports * the same query methods with promise pipelining. * * @example * ```typescript * await db.transaction(async tx => { * await tx.execute('INSERT INTO users (name) VALUES ($1)', ['Alice']) * await tx.execute('INSERT INTO audit_log (action) VALUES ($1)', ['user_created']) * }) * ``` * * @example Promise pipelining in transactions * ```typescript * await db.transaction(tx => { * const balance = tx.query('SELECT balance FROM accounts WHERE id = $1', [fromId]) * return balance.map(async rows => { * if (rows[0].balance < amount) throw new Error('Insufficient funds') * await tx.execute('UPDATE accounts SET balance = balance - $1 WHERE id = $2', [amount, fromId]) * await tx.execute('UPDATE accounts SET balance = balance + $1 WHERE id = $2', [amount, toId]) * }) * }) * ``` */ async transaction( fn: (tx: TransactionClient) => Promise | RpcPromise, options?: TransactionOptions ): Promise { // Start collecting queries for batching this.transport.startCollecting() // Build BEGIN statement const beginParts = ['BEGIN'] if (options?.isolationLevel) { // Validate isolation level to prevent SQL injection const validLevels = ['read uncommitted', 'read committed', 'repeatable read', 'serializable'] const normalized = options.isolationLevel.toLowerCase() if (!validLevels.includes(normalized)) { throw new Error(`Invalid isolation level: "${options.isolationLevel}"`) } beginParts.push(`ISOLATION LEVEL ${normalized.toUpperCase()}`) } if (options?.readOnly) { beginParts.push('READ ONLY') } try { // Execute BEGIN await this.transport.query(beginParts.join(' ')) // Create transaction client const txClient: TransactionClient = { query: (sql: string, params?: unknown[]) => this.query(sql, params), queryOne: (sql: string, params?: unknown[]) => this.queryOne(sql, params), queryScalar: (sql: string, params?: unknown[]) => this.queryScalar(sql, params), execute: (sql: string, params?: unknown[]) => this.execute(sql, params), } // Execute transaction callback const result = fn(txClient) const finalResult = result instanceof RpcPromise ? await result.execute() : await result // Stop collecting and execute any remaining batched queries await this.transport.stopCollecting() // Commit await this.transport.query('COMMIT') return finalResult } catch (error) { // Stop collecting this.transport.stopCollecting().catch(() => {}) // Rollback await this.transport.query('ROLLBACK').catch(() => {}) throw error } } /** * Execute queries within a batched transaction * * All queries are collected and executed in a single round trip, * wrapped in BEGIN/COMMIT. * * @example * ```typescript * const results = await db.batchTransaction([ * { sql: 'INSERT INTO users (name) VALUES ($1)', params: ['Alice'] }, * { sql: 'INSERT INTO audit_log (action) VALUES ($1)', params: ['user_created'] }, * ]) * ``` */ async batchTransaction( queries: Array<{ sql: string; params?: unknown[] }>, options?: TransactionOptions ): Promise { const results = await this.transport.batchTransaction(queries, options) return results.map(r => r.rows) } /** * Health check - verify database is responsive */ ping(): RpcPromise<{ ok: true; durationMs: number }> { return createRpcPromise(async () => { const start = performance.now() await this.transport.query('SELECT 1') return { ok: true as const, durationMs: performance.now() - start } }) } /** * Get database version */ version(): RpcPromise { return createRpcPromise(async () => { const rows = await this.query<{ version: string }>('SELECT version()') return rows[0]?.version ?? 'unknown' }) } /** * List all tables in the public schema */ listTables(): RpcPromise { return createRpcPromise(async () => { const rows = await this.query<{ tablename: string }>( `SELECT tablename FROM pg_tables WHERE schemaname = 'public' ORDER BY tablename` ) return rows.map(r => r.tablename) }) } /** * Get table schema information */ describeTable(tableName: string): RpcPromise> { return createRpcPromise(async () => { const rows = await this.query<{ column_name: string data_type: string is_nullable: string column_default: string | null }>( `SELECT column_name, data_type, is_nullable, column_default FROM information_schema.columns WHERE table_schema = 'public' AND table_name = $1 ORDER BY ordinal_position`, [tableName] ) return rows.map(r => ({ column_name: r.column_name, data_type: r.data_type, is_nullable: r.is_nullable === 'YES', column_default: r.column_default, })) }) } /** * Close the connection */ async close(): Promise { await this.transport.close() } /** * Check if connected */ isConnected(): boolean { return this.transport.isConnected() } } /** * Create a capnweb-enabled PostgreSQL client * * @example * ```typescript * import { createClient } from 'postgres.do/rpc' * * const db = createClient('postgres.do/my-database') * * // Query with magic map * const usersWithOrders = await db.query('SELECT * FROM users') * .map(users => Promise.all( * users.map(u => db.query('SELECT * FROM orders WHERE user_id = $1', [u.id])) * )) * ``` */ export function createClient(urlOrConfig?: string | RpcClientConfig): RpcClient { if (typeof urlOrConfig === 'string') { return new RpcClient({ url: urlOrConfig }) } return new RpcClient(urlOrConfig) } /** * Re-export utilities */ export { RpcPromise, createRpcPromise, magicMap, batchExecute } export type { RpcQueryResult }