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<AppliedMigration>(
"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;
}
}