import type pg from "pg"; export type PgD1Database = { prepare(sql: string): PgD1PreparedStatement; batch<_T = any>(statements: PgD1PreparedStatement[]): Promise; query(text: string, params?: unknown[]): Promise>; }; export interface PgD1PreparedStatement { bind(...params: unknown[]): PgD1PreparedStatement; all(): Promise<{ results: T[] }>; first(): Promise; run(): Promise<{ meta: { changes: number } }>; } export function createPgD1(pool: pg.Pool): PgD1Database { return { prepare(sql: string) { return createStatement(pool, sql); }, async batch<_T = any>(statements: PgD1PreparedStatement[]): Promise { const client = await pool.connect(); try { await client.query("BEGIN"); const results = []; for (const stmt of statements) { const s = stmt as StatementImpl; const { queryText, queryParams } = convertSqlForPostgres(s.sql, s.params); const res = await client.query(queryText, queryParams); results.push({ results: res.rows, meta: { changes: res.rowCount ?? 0 } }); } await client.query("COMMIT"); return results; } catch (err) { await client.query("ROLLBACK"); throw err; } finally { client.release(); } }, async query(text: string, params?: unknown[]): Promise> { const { queryText, queryParams } = convertSqlForPostgres(text, params || []); return pool.query(queryText, queryParams); }, }; } class StatementImpl implements PgD1PreparedStatement { pool: pg.Pool; sql: string; params: unknown[] = []; constructor(pool: pg.Pool, sql: string) { this.pool = pool; this.sql = sql; } bind(...params: unknown[]): PgD1PreparedStatement { this.params = params; return this; } async all(): Promise<{ results: T[] }> { const { queryText, queryParams } = convertSqlForPostgres(this.sql, this.params); const res = await this.pool.query(queryText, queryParams); return { results: res.rows }; } async first(): Promise { const { results } = await this.all(); return results[0] || null; } async run(): Promise<{ meta: { changes: number } }> { const { queryText, queryParams } = convertSqlForPostgres(this.sql, this.params); const res = await this.pool.query(queryText, queryParams); return { meta: { changes: res.rowCount ?? 0 } }; } } function createStatement(pool: pg.Pool, sql: string): PgD1PreparedStatement { return new StatementImpl(pool, sql); } function convertSqlForPostgres(sql: string, params: unknown[]): { queryText: string; queryParams: unknown[] } { let converted = sql; // 1. Convert SQLite "INSERT OR IGNORE" to PostgreSQL "INSERT ... ON CONFLICT DO NOTHING" if (/\bINSERT\s+OR\s+IGNORE\s+INTO\b/i.test(converted)) { converted = converted.replace(/\bINSERT\s+OR\s+IGNORE\s+INTO\b/gi, "INSERT INTO"); if (!/\bON\s+CONFLICT\b/i.test(converted)) { converted += " ON CONFLICT DO NOTHING"; } } // 2. Convert SQLite "INSERT OR REPLACE" to PostgreSQL "INSERT ... ON CONFLICT (...) DO UPDATE" if (/\bINSERT\s+OR\s+REPLACE\s+INTO\b/i.test(converted)) { converted = converted.replace(/\bINSERT\s+OR\s+REPLACE\s+INTO\b/gi, "INSERT INTO"); } // 3. Convert SQLite datetime('now') or strftime() to PostgreSQL NOW() converted = converted.replace(/\bdatetime\('now'\)/gi, "NOW()"); converted = converted.replace(/\bdatetime\('now',\s*'localtime'\)/gi, "NOW()"); converted = converted.replace(/\bstrftime\([^)]+\)/gi, "NOW()"); // 4. Convert SQLite placeholders "?" to PostgreSQL "$1, $2, $3..." let paramIndex = 1; converted = converted.replace(/\?/g, () => { return `$${paramIndex++}`; }); return { queryText: converted, queryParams: params, }; }