sqlite.bun.ts 6.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183
  1. import { Database } from "bun:sqlite"
  2. import { drizzle } from "drizzle-orm/bun-sqlite"
  3. import * as Context from "effect/Context"
  4. import * as Effect from "effect/Effect"
  5. import * as Fiber from "effect/Fiber"
  6. import { identity } from "effect/Function"
  7. import * as Layer from "effect/Layer"
  8. import * as Scope from "effect/Scope"
  9. import * as Semaphore from "effect/Semaphore"
  10. import * as Stream from "effect/Stream"
  11. import * as Reactivity from "effect/unstable/reactivity/Reactivity"
  12. import * as Client from "effect/unstable/sql/SqlClient"
  13. import type { Connection } from "effect/unstable/sql/SqlConnection"
  14. import { classifySqliteError, SqlError } from "effect/unstable/sql/SqlError"
  15. import * as Statement from "effect/unstable/sql/Statement"
  16. import { Sqlite } from "./sqlite"
  17. const ATTR_DB_SYSTEM_NAME = "db.system.name"
  18. const TypeId = "~@opencode-ai/core/database/SqliteBun" as const
  19. type TypeId = typeof TypeId
  20. interface SqliteClient extends Client.SqlClient {
  21. readonly [TypeId]: TypeId
  22. readonly config: Config
  23. readonly export: Effect.Effect<Uint8Array, SqlError>
  24. readonly loadExtension: (path: string) => Effect.Effect<void, SqlError>
  25. readonly updateValues: never
  26. }
  27. interface Config {
  28. readonly filename: string
  29. readonly readonly?: boolean
  30. readonly create?: boolean
  31. readonly readwrite?: boolean
  32. readonly disableWAL?: boolean
  33. readonly spanAttributes?: Record<string, unknown>
  34. readonly transformResultNames?: (str: string) => string
  35. readonly transformQueryNames?: (str: string) => string
  36. }
  37. interface SqliteConnection extends Connection {
  38. readonly export: Effect.Effect<Uint8Array, SqlError>
  39. readonly loadExtension: (path: string) => Effect.Effect<void, SqlError>
  40. }
  41. const make = (options: Config) =>
  42. Effect.gen(function* () {
  43. const native = (yield* Sqlite.Native) as Database
  44. const compiler = Statement.makeCompilerSqlite(options.transformQueryNames)
  45. const transformRows = options.transformResultNames
  46. ? Statement.defaultTransforms(options.transformResultNames).array
  47. : undefined
  48. const run = (query: string, params: ReadonlyArray<unknown> = []) =>
  49. Effect.withFiber<Array<Record<string, unknown>>, SqlError>((fiber) => {
  50. const statement = native.query(query)
  51. // @ts-ignore bun-types missing safeIntegers method, fixed in https://github.com/oven-sh/bun/pull/26627
  52. statement.safeIntegers(Context.get(fiber.context, Client.SafeIntegers))
  53. try {
  54. return Effect.succeed((statement.all(...(params as any)) ?? []) as Array<Record<string, unknown>>)
  55. } catch (cause) {
  56. return Effect.fail(
  57. new SqlError({
  58. reason: classifySqliteError(cause, { message: "Failed to execute statement", operation: "execute" }),
  59. }),
  60. )
  61. }
  62. })
  63. const runValues = (query: string, params: ReadonlyArray<unknown> = []) =>
  64. Effect.withFiber<Array<unknown[]>, SqlError>((fiber) => {
  65. const statement = native.query(query)
  66. // @ts-ignore bun-types missing safeIntegers method, fixed in https://github.com/oven-sh/bun/pull/26627
  67. statement.safeIntegers(Context.get(fiber.context, Client.SafeIntegers))
  68. try {
  69. return Effect.succeed((statement.values(...(params as any)) ?? []) as Array<unknown[]>)
  70. } catch (cause) {
  71. return Effect.fail(
  72. new SqlError({
  73. reason: classifySqliteError(cause, { message: "Failed to execute statement", operation: "execute" }),
  74. }),
  75. )
  76. }
  77. })
  78. const connection = identity<SqliteConnection>({
  79. execute(query, params, transformRows) {
  80. return transformRows ? Effect.map(run(query, params), transformRows) : run(query, params)
  81. },
  82. executeRaw(query, params) {
  83. return run(query, params)
  84. },
  85. executeValues(query, params) {
  86. return runValues(query, params)
  87. },
  88. executeUnprepared(query, params, transformRows) {
  89. return this.execute(query, params, transformRows)
  90. },
  91. executeStream() {
  92. return Stream.die("executeStream not implemented")
  93. },
  94. export: Effect.try({
  95. try: () => native.serialize(),
  96. catch: (cause) =>
  97. new SqlError({
  98. reason: classifySqliteError(cause, { message: "Failed to export database", operation: "export" }),
  99. }),
  100. }),
  101. loadExtension: (path) =>
  102. Effect.try({
  103. try: () => native.loadExtension(path),
  104. catch: (cause) =>
  105. new SqlError({
  106. reason: classifySqliteError(cause, { message: "Failed to load extension", operation: "loadExtension" }),
  107. }),
  108. }),
  109. })
  110. const semaphore = yield* Semaphore.make(1)
  111. const acquirer = semaphore.withPermits(1)(Effect.succeed(connection))
  112. const transactionAcquirer = Effect.uninterruptibleMask((restore) => {
  113. const fiber = Fiber.getCurrent()!
  114. const scope = Context.getUnsafe(fiber.context, Scope.Scope)
  115. return Effect.as(
  116. Effect.tap(restore(semaphore.take(1)), () => Scope.addFinalizer(scope, semaphore.release(1))),
  117. connection,
  118. )
  119. })
  120. const client = Object.assign(
  121. (yield* Client.make({
  122. acquirer,
  123. compiler,
  124. transactionAcquirer,
  125. spanAttributes: [
  126. ...(options.spanAttributes ? Object.entries(options.spanAttributes) : []),
  127. [ATTR_DB_SYSTEM_NAME, "sqlite"],
  128. ],
  129. transformRows,
  130. })) as SqliteClient,
  131. {
  132. [TypeId]: TypeId,
  133. config: options,
  134. export: Effect.flatMap(acquirer, (_) => _.export),
  135. loadExtension: (path: string) => Effect.flatMap(acquirer, (_) => _.loadExtension(path)),
  136. },
  137. )
  138. return client
  139. })
  140. const nativeLayer = (config: Config) =>
  141. Layer.effect(
  142. Sqlite.Native,
  143. Effect.gen(function* () {
  144. const native = new Database(config.filename, {
  145. readonly: config.readonly,
  146. readwrite: config.readwrite ?? true,
  147. create: config.create ?? true,
  148. })
  149. yield* Effect.addFinalizer(() => Effect.sync(() => native.close()))
  150. if (config.disableWAL !== true) native.run("PRAGMA journal_mode = WAL;")
  151. return native
  152. }),
  153. )
  154. const sqliteLayer = (config: Config) => Layer.effect(Client.SqlClient, make(config))
  155. const drizzleLayer = Layer.effect(
  156. Sqlite.Drizzle,
  157. Effect.gen(function* () {
  158. return drizzle({ client: (yield* Sqlite.Native) as Database })
  159. }),
  160. )
  161. export const layer = (config: Config) => {
  162. const native = nativeLayer(config)
  163. return Layer.merge(native, Layer.merge(sqliteLayer(config), drizzleLayer).pipe(Layer.provide(native))).pipe(
  164. Layer.provide(Reactivity.layer),
  165. )
  166. }