import { readdir, readFile } from "node:fs/promises"; import { join } from "node:path"; import type { Queryable } from "@/lib/db/types"; type AppliedMigration = { name: string }; export async function applyMigrations( database: Queryable, migrationsDirectory: string, ) { await database.query(` CREATE TABLE IF NOT EXISTS schema_migrations ( name text PRIMARY KEY, applied_at timestamptz NOT NULL DEFAULT now() ) `); const migrationFiles = (await readdir(migrationsDirectory)) .filter((file) => /^\d+_.+\.sql$/.test(file)) .sort(); const applied = await database.query( "SELECT name FROM schema_migrations", ); const appliedNames = new Set(applied.rows.map((migration) => migration.name)); const pending = migrationFiles.filter((file) => !appliedNames.has(file)); if (pending.length === 0) { return; } await database.query("BEGIN"); try { for (const migrationName of pending) { const sql = await readFile(join(migrationsDirectory, migrationName), "utf8"); await database.query(sql); await database.query( "INSERT INTO schema_migrations (name) VALUES ($1)", [migrationName], ); } await database.query("COMMIT"); } catch (error) { await database.query("ROLLBACK"); throw error; } }