effect-sqlite.ts 3.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102
  1. /* oxlint-disable */
  2. import * as Effect from "effect/Effect"
  3. import type { SqlError } from "effect/unstable/sql/SqlError"
  4. import { EffectDrizzleError } from "drizzle-orm/effect-core/errors"
  5. import type { QueryEffectHKTBase } from "drizzle-orm/effect-core/query-effect"
  6. import type { MigrationMeta } from "drizzle-orm/migrator"
  7. import { sql } from "drizzle-orm/sql/sql"
  8. import type { SQLiteEffectSession } from "../sqlite-core/effect/session"
  9. import {
  10. buildSQLiteMigrationBackfillStatements,
  11. prepareSQLiteMigrationBackfill,
  12. type SQLiteMigrationTableRow,
  13. } from "./sqlite"
  14. import { GET_VERSION_FOR, MIGRATIONS_TABLE_VERSIONS, type UpgradeResult } from "./utils"
  15. const migrationUpgradeError = (cause: unknown) =>
  16. new EffectDrizzleError({
  17. message:
  18. typeof cause === "object" && cause !== null && "message" in cause && typeof cause.message === "string"
  19. ? cause.message
  20. : String(cause),
  21. cause,
  22. })
  23. export const upgradeIfNeeded: <TEffectHKT extends QueryEffectHKTBase>(
  24. migrationsTable: string,
  25. session: SQLiteEffectSession<TEffectHKT>,
  26. localMigrations: MigrationMeta[],
  27. ) => Effect.Effect<UpgradeResult, EffectDrizzleError | TEffectHKT["error"] | SqlError, TEffectHKT["context"]> =
  28. Effect.fn("upgradeIfNeeded")(function* <TEffectHKT extends QueryEffectHKTBase>(
  29. migrationsTable: string,
  30. session: SQLiteEffectSession<TEffectHKT>,
  31. localMigrations: MigrationMeta[],
  32. ) {
  33. const tableExists = yield* session.all(
  34. sql`SELECT 1 FROM sqlite_master WHERE type = 'table' AND name = ${migrationsTable}`,
  35. )
  36. if (tableExists.length === 0) {
  37. return { newDb: true }
  38. }
  39. const rows = yield* session.all<{ column_name: string }>(
  40. sql`SELECT name as column_name FROM pragma_table_info(${migrationsTable})`,
  41. )
  42. const version = GET_VERSION_FOR.sqlite(rows.map((r) => r.column_name))
  43. for (let v = version; v < MIGRATIONS_TABLE_VERSIONS.sqlite; v++) {
  44. const upgradeFn = upgradeFunctions[v]
  45. if (!upgradeFn) {
  46. return yield* new EffectDrizzleError({
  47. message: `No upgrade path from migration table version ${v} to ${v + 1}`,
  48. cause: { version: v },
  49. })
  50. }
  51. yield* upgradeFn(migrationsTable, session, localMigrations)
  52. }
  53. return { newDb: false }
  54. })
  55. const upgradeFunctions: Record<
  56. number,
  57. <TEffectHKT extends QueryEffectHKTBase>(
  58. migrationsTable: string,
  59. session: SQLiteEffectSession<TEffectHKT>,
  60. localMigrations: MigrationMeta[],
  61. ) => Effect.Effect<void, EffectDrizzleError | TEffectHKT["error"] | SqlError, TEffectHKT["context"]>
  62. > = {
  63. 0: upgradeFromV0,
  64. }
  65. function upgradeFromV0<TEffectHKT extends QueryEffectHKTBase>(
  66. migrationsTable: string,
  67. session: SQLiteEffectSession<TEffectHKT>,
  68. localMigrations: MigrationMeta[],
  69. ): Effect.Effect<void, EffectDrizzleError | TEffectHKT["error"] | SqlError, TEffectHKT["context"]> {
  70. return Effect.gen(function* () {
  71. const table = sql`${sql.identifier(migrationsTable)}`
  72. const dbRows = yield* session.all<SQLiteMigrationTableRow>(
  73. sql`SELECT id, hash, created_at FROM ${table} ORDER BY id ASC`,
  74. )
  75. const statements = yield* Effect.try({
  76. try: () =>
  77. buildSQLiteMigrationBackfillStatements(
  78. migrationsTable,
  79. prepareSQLiteMigrationBackfill(dbRows, localMigrations),
  80. ),
  81. catch: migrationUpgradeError,
  82. })
  83. yield* session.transaction((tx) =>
  84. Effect.gen(function* () {
  85. for (const statement of statements) {
  86. yield* tx.run(statement)
  87. }
  88. }),
  89. )
  90. })
  91. }