/
Alex_Ural
/
opencode
Обзор
Документация
Войти
/
Alex_Ural
/
opencode
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
Аналитика
Безопасность
dev
packages/effect-drizzle-sqlite/src/effect-sqlite/session.ts
218 строк
9 KB
Kit Langton
fix(effect-drizzle-sqlite): preserve transaction begin errors (#30448)
02 июн 2026, 21:36
Не верифицирован
02 июн 2026, 21:36
4002b85
Код
Авторство
О чём код?
/* oxlint-disable */ import * as Context from "effect/Context" import * as Effect from "effect/Effect" import * as Exit from "effect/Exit" import * as Scope from "effect/Scope" import type { SqlClient } from "effect/unstable/sql/SqlClient" import type { SqlError } from "effect/unstable/sql/SqlError" import type { EffectCacheShape } from "drizzle-orm/cache/core/cache-effect" import type { WithCacheConfig } from "drizzle-orm/cache/core/types" import type { EffectDrizzleQueryError } from "drizzle-orm/effect-core/errors" import type { EffectLoggerShape } from "drizzle-orm/effect-core/logger" import type { QueryEffectHKTBase } from "drizzle-orm/effect-core/query-effect" import { entityKind } from "drizzle-orm/entity" import type { AnyRelations } from "drizzle-orm/relations" import type { RelationalQueryMapperConfig } from "drizzle-orm/relations" import type { Query } from "drizzle-orm/sql/sql" import type { SQLiteAsyncDialect } from "drizzle-orm/sqlite-core/dialect" import { SQLiteEffectPreparedQuery, SQLiteEffectSession, SQLiteEffectTransaction } from "../sqlite-core/effect/session" import type { SelectedFieldsOrdered } from "drizzle-orm/sqlite-core/query-builders/select.types" import type { PreparedQueryConfig, SQLiteExecuteMethod, SQLiteTransactionConfig } from "drizzle-orm/sqlite-core/session" export interface EffectSQLiteQueryEffectHKT extends QueryEffectHKTBase { readonly error: EffectDrizzleQueryError readonly context: never } export type EffectSQLiteRunResult = readonly never[] export interface EffectSQLiteSessionOptions { logger: EffectLoggerShape cache: EffectCacheShape useJitMappers?: boolean } export class EffectSQLiteSession<TRelations extends AnyRelations> extends SQLiteEffectSession< EffectSQLiteQueryEffectHKT, EffectSQLiteRunResult, TRelations > { static override readonly [entityKind]: string = "EffectSQLiteSession" constructor( private client: SqlClient, dialect: SQLiteAsyncDialect, protected relations: TRelations, private options: EffectSQLiteSessionOptions, ) { super(dialect) } override prepareQuery<T extends PreparedQueryConfig = PreparedQueryConfig>( query: Query, fields: SelectedFieldsOrdered | undefined, executeMethod: SQLiteExecuteMethod, customResultMapper?: (rows: unknown[][], mapColumnValue?: (value: unknown) => unknown) => unknown, queryMetadata?: { type: "select" | "update" | "delete" | "insert" tables: string[] }, cacheConfig?: WithCacheConfig, ): SQLiteEffectPreparedQuery<T, EffectSQLiteQueryEffectHKT> { return new SQLiteEffectPreparedQuery<T, EffectSQLiteQueryEffectHKT>( (params, method) => this.execute(query, params, method), query, this.options.logger, this.options.cache, queryMetadata, cacheConfig, fields, executeMethod, this.options.useJitMappers, customResultMapper, undefined, undefined, this.isInTransaction(), ) } override prepareRelationalQuery<T extends PreparedQueryConfig = PreparedQueryConfig>( query: Query, fields: SelectedFieldsOrdered | undefined, executeMethod: SQLiteExecuteMethod, customResultMapper: (rows: Record<string, unknown>[], mapColumnValue?: (value: unknown) => unknown) => unknown, config: RelationalQueryMapperConfig, ): SQLiteEffectPreparedQuery<T, EffectSQLiteQueryEffectHKT, true> { return new SQLiteEffectPreparedQuery<T, EffectSQLiteQueryEffectHKT, true>( (params, method) => this.execute(query, params, method), query, this.options.logger, this.options.cache, undefined, undefined, fields, executeMethod, this.options.useJitMappers, customResultMapper, true, config, this.isInTransaction(), ) } private execute(query: Query, params: unknown[], method: SQLiteExecuteMethod | "values") { const statement = this.client.unsafe(query.sql, params) if (method === "values") return statement.values if (method === "get") return statement.withoutTransform.pipe(Effect.map((rows) => rows[0])) return statement.withoutTransform } private isInTransaction() { return Effect.serviceOption(this.client.transactionService).pipe(Effect.map((option) => option._tag === "Some")) } private executeTransactionStatement(connection: Effect.Success<SqlClient["reserve"]>, query: string) { return connection.executeUnprepared(query, [], undefined).pipe(Effect.asVoid) } private withTransaction<A, E, R>(effect: Effect.Effect<A, E, R>, config: SQLiteTransactionConfig | undefined) { return Effect.uninterruptibleMask((restore) => Effect.withFiber<A, E | SqlError, R>((fiber) => { const services = fiber.context const connectionOption = Context.getOption(services, this.client.transactionService) const connection: Effect.Effect< readonly [Scope.Closeable | undefined, Effect.Success<SqlClient["reserve"]>], SqlError > = connectionOption._tag === "Some" ? Effect.succeed([undefined, connectionOption.value[0]] as const) : Scope.make().pipe( Effect.flatMap((scope) => Scope.provide(this.client.reserve, scope).pipe( Effect.map((connection) => [scope, connection] as const), Effect.catch((error) => Scope.close(scope, Exit.fail(error)).pipe(Effect.andThen(Effect.fail(error))), ), ), ), ) const id = connectionOption._tag === "Some" ? connectionOption.value[1] + 1 : 0 return connection.pipe( Effect.flatMap(([scope, connection]) => { const transaction = this.executeTransactionStatement( connection, id === 0 ? `begin ${config?.behavior ?? "deferred"}` : `savepoint effect_sql_${id}`, ).pipe( Effect.flatMap(() => Effect.provideContext( restore(effect), Context.add(services, this.client.transactionService, [connection, id]), ).pipe( Effect.exit, Effect.flatMap((exit) => { const finalize = Exit.isSuccess(exit) ? id === 0 ? this.executeTransactionStatement(connection, "commit").pipe( // SQLite keeps the transaction open after deferred constraint commit failures. Effect.catch((error) => this.executeTransactionStatement(connection, "rollback").pipe( Effect.catch(() => Effect.void), Effect.andThen(Effect.fail(error)), ), ), ) : this.executeTransactionStatement(connection, `release savepoint effect_sql_${id}`) : id === 0 ? this.executeTransactionStatement(connection, "rollback") : this.executeTransactionStatement(connection, `rollback to savepoint effect_sql_${id}`).pipe( Effect.andThen( this.executeTransactionStatement(connection, `release savepoint effect_sql_${id}`), ), ) return finalize.pipe(Effect.flatMap(() => exit)) }), ), ), ) return scope === undefined ? transaction : transaction.pipe(Effect.onExit((exit) => Scope.close(scope, exit))) }), ) }), ) } override transaction<A, E, R>( transaction: (tx: EffectSQLiteTransaction<TRelations>) => Effect.Effect<A, E, R>, config?: SQLiteTransactionConfig, ): Effect.Effect<A, E | SqlError, R> { const { dialect, relations } = this return this.withTransaction( Effect.gen({ self: this }, function* () { const tx = new EffectSQLiteTransaction<TRelations>(dialect, this, relations) return yield* transaction(tx) }), config, ) } } export class EffectSQLiteTransaction<TRelations extends AnyRelations> extends SQLiteEffectTransaction< EffectSQLiteQueryEffectHKT, EffectSQLiteRunResult, TRelations > { static override readonly [entityKind]: string = "EffectSQLiteTransaction" override transaction: <A, E, R>( transaction: ( tx: SQLiteEffectTransaction<EffectSQLiteQueryEffectHKT, EffectSQLiteRunResult, TRelations>, ) => Effect.Effect<A, E, R>, ) => Effect.Effect<A, SqlError | E, R> = (tx) => this.session.transaction(tx) }