index.ts 6.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168
  1. export * as NodeSqliteClient from "./index"
  2. import { DatabaseSync, type SQLInputValue } from "node:sqlite"
  3. import { identity } from "effect/Function"
  4. import * as Context from "effect/Context"
  5. import * as Effect from "effect/Effect"
  6. import * as Fiber from "effect/Fiber"
  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. const ATTR_DB_SYSTEM_NAME = "db.system.name"
  17. export const TypeId: TypeId = "~@opencode-ai/effect-sqlite-node/NodeSqliteClient"
  18. export type TypeId = "~@opencode-ai/effect-sqlite-node/NodeSqliteClient"
  19. export interface SqliteClient extends Client.SqlClient {
  20. readonly [TypeId]: TypeId
  21. readonly config: SqliteClientConfig
  22. readonly loadExtension: (path: string) => Effect.Effect<void, SqlError>
  23. readonly updateValues: never
  24. }
  25. export const SqliteClient = Context.Service<SqliteClient>("@opencode-ai/effect-sqlite-node/NodeSqliteClient")
  26. export interface SqliteClientConfig {
  27. readonly filename: string
  28. readonly readonly?: boolean | undefined
  29. readonly create?: boolean | undefined
  30. readonly readwrite?: boolean | undefined
  31. readonly disableWAL?: boolean | undefined
  32. readonly timeout?: number | undefined
  33. readonly allowExtension?: boolean | undefined
  34. readonly spanAttributes?: Record<string, unknown> | undefined
  35. readonly transformResultNames?: ((str: string) => string) | undefined
  36. readonly transformQueryNames?: ((str: string) => string) | undefined
  37. }
  38. interface SqliteConnection extends Connection {
  39. readonly loadExtension: (path: string) => Effect.Effect<void, SqlError>
  40. }
  41. export const make = (
  42. options: SqliteClientConfig,
  43. ): Effect.Effect<SqliteClient, never, Scope.Scope | Reactivity.Reactivity> =>
  44. Effect.gen(function* () {
  45. const compiler = Statement.makeCompilerSqlite(options.transformQueryNames)
  46. const transformRows = options.transformResultNames
  47. ? Statement.defaultTransforms(options.transformResultNames).array
  48. : undefined
  49. const makeConnection = Effect.gen(function* () {
  50. const db = new DatabaseSync(options.filename, {
  51. readOnly: options.readonly,
  52. timeout: options.timeout,
  53. allowExtension: options.allowExtension,
  54. enableForeignKeyConstraints: true,
  55. open: true,
  56. })
  57. yield* Effect.addFinalizer(() => Effect.sync(() => db.close()))
  58. if (options.disableWAL !== true && options.readonly !== true) {
  59. db.exec("PRAGMA journal_mode = WAL;")
  60. }
  61. const run = (sql: string, params: ReadonlyArray<unknown> = []) =>
  62. Effect.withFiber<Array<Record<string, unknown>>, SqlError>((fiber) => {
  63. const statement = db.prepare(sql)
  64. statement.setReadBigInts(Context.get(fiber.context, Client.SafeIntegers))
  65. try {
  66. return Effect.succeed(statement.all(...(params as SQLInputValue[])) as Array<Record<string, unknown>>)
  67. } catch (cause) {
  68. return Effect.fail(
  69. new SqlError({
  70. reason: classifySqliteError(cause, { message: "Failed to execute statement", operation: "execute" }),
  71. }),
  72. )
  73. }
  74. })
  75. const runValues = (sql: string, params: ReadonlyArray<unknown> = []) =>
  76. Effect.withFiber<ReadonlyArray<ReadonlyArray<unknown>>, SqlError>((fiber) => {
  77. const statement = db.prepare(sql)
  78. statement.setReadBigInts(Context.get(fiber.context, Client.SafeIntegers))
  79. statement.setReturnArrays(true)
  80. try {
  81. return Effect.succeed(
  82. statement.all(...(params as SQLInputValue[])) as unknown as ReadonlyArray<ReadonlyArray<unknown>>,
  83. )
  84. } catch (cause) {
  85. return Effect.fail(
  86. new SqlError({
  87. reason: classifySqliteError(cause, { message: "Failed to execute statement", operation: "execute" }),
  88. }),
  89. )
  90. }
  91. })
  92. return identity<SqliteConnection>({
  93. execute(sql, params, transformRows) {
  94. return transformRows ? Effect.map(run(sql, params), transformRows) : run(sql, params)
  95. },
  96. executeRaw(sql, params) {
  97. return run(sql, params)
  98. },
  99. executeValues(sql, params) {
  100. return runValues(sql, params)
  101. },
  102. executeUnprepared(sql, params, transformRows) {
  103. return this.execute(sql, params, transformRows)
  104. },
  105. executeStream() {
  106. return Stream.die("executeStream not implemented")
  107. },
  108. loadExtension: (path) =>
  109. Effect.try({
  110. try: () => db.loadExtension(path),
  111. catch: (cause) =>
  112. new SqlError({
  113. reason: classifySqliteError(cause, { message: "Failed to load extension", operation: "loadExtension" }),
  114. }),
  115. }),
  116. })
  117. })
  118. const semaphore = yield* Semaphore.make(1)
  119. const connection = yield* makeConnection
  120. const acquirer = semaphore.withPermits(1)(Effect.succeed(connection))
  121. const transactionAcquirer = Effect.uninterruptibleMask((restore) => {
  122. const fiber = Fiber.getCurrent()!
  123. const scope = Context.getUnsafe(fiber.context, Scope.Scope)
  124. return Effect.as(
  125. Effect.tap(restore(semaphore.take(1)), () => Scope.addFinalizer(scope, semaphore.release(1))),
  126. connection,
  127. )
  128. })
  129. return Object.assign(
  130. (yield* Client.make({
  131. acquirer,
  132. compiler,
  133. transactionAcquirer,
  134. spanAttributes: [
  135. ...(options.spanAttributes ? Object.entries(options.spanAttributes) : []),
  136. [ATTR_DB_SYSTEM_NAME, "sqlite"],
  137. ],
  138. transformRows,
  139. })) as SqliteClient,
  140. {
  141. [TypeId]: TypeId as TypeId,
  142. config: options,
  143. loadExtension: (path: string) => Effect.flatMap(acquirer, (_) => _.loadExtension(path)),
  144. },
  145. )
  146. })
  147. export const layer = (config: SqliteClientConfig): Layer.Layer<SqliteClient | Client.SqlClient> =>
  148. Layer.effectContext(
  149. Effect.map(make(config), (client) =>
  150. Context.make(SqliteClient, client).pipe(Context.add(Client.SqlClient, client)),
  151. ),
  152. ).pipe(Layer.provide(Reactivity.layer))