aegida-console / lib / db / migrate.ts
migrate.ts
Raw
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;
  }
}