fix(effect-drizzle-sqlite): support pipeable transactions

This commit is contained in:
Kit Langton 2026-04-27 22:18:44 -04:00
commit e4ae265d8f
2 changed files with 68 additions and 39 deletions

View file

@ -31,10 +31,10 @@ export type EffectSQLiteDatabase<
TRelations extends AnyRelations = EmptyRelations, TRelations extends AnyRelations = EmptyRelations,
> = SQLiteBunDatabase<TSchema, TRelations> & { > = SQLiteBunDatabase<TSchema, TRelations> & {
readonly $client: Database readonly $client: Database
readonly withTransaction: <A, E>( readonly withTransaction: <A, E, R>(
transaction: (tx: SQLiteTransaction<"sync", void, TSchema, TRelations>) => Effect.Effect<A, E>, effect: Effect.Effect<A, E, R>,
config?: SQLiteTransactionConfig, config?: SQLiteTransactionConfig,
) => Effect.Effect<A, E> ) => Effect.Effect<A, E, R>
} }
export type MakeConfig< export type MakeConfig<
@ -176,33 +176,48 @@ const attachTransaction = <
TSchema extends Record<string, unknown> = Record<string, never>, TSchema extends Record<string, unknown> = Record<string, never>,
TRelations extends AnyRelations = EmptyRelations, TRelations extends AnyRelations = EmptyRelations,
>(db: SQLiteBunDatabase<TSchema, TRelations> & { readonly $client: Database }): EffectSQLiteDatabase<TSchema, TRelations> => { >(db: SQLiteBunDatabase<TSchema, TRelations> & { readonly $client: Database }): EffectSQLiteDatabase<TSchema, TRelations> => {
const runTransaction = db.transaction.bind(db) as ( const txStack: Array<SQLiteTransaction<"sync", void, TSchema, TRelations>> = []
transaction: (tx: SQLiteTransaction<"sync", void, TSchema, TRelations>) => unknown, const current = () => txStack.at(-1) ?? db
config?: SQLiteTransactionConfig, const runTransaction = (target: SQLiteBunDatabase<TSchema, TRelations> | SQLiteTransaction<"sync", void, TSchema, TRelations>) =>
) => unknown target.transaction.bind(target) as (
transaction: (tx: SQLiteTransaction<"sync", void, TSchema, TRelations>) => unknown,
return Object.assign(db, {
withTransaction: <A, E>(
transaction: (tx: SQLiteTransaction<"sync", void, TSchema, TRelations>) => Effect.Effect<A, E>,
config?: SQLiteTransactionConfig, config?: SQLiteTransactionConfig,
) => ) => unknown
Effect.sync(
() => const withTransaction = <A, E, R>(
runTransaction( effect: Effect.Effect<A, E, R>,
(tx) => config?: SQLiteTransactionConfig,
Exit.match(Effect.runSyncExit(transaction(tx)), { ): Effect.Effect<A, E, R> =>
onSuccess: (value) => value, Effect.context<R>().pipe(
onFailure: (cause) => { Effect.flatMap((context) =>
throw new TransactionFailure(cause) Effect.sync(
}, () =>
}), runTransaction(current())((tx) => {
config, txStack.push(tx)
) as A, try {
).pipe( const exit = Effect.runSyncExit(Effect.provideContext(effect, context))
Effect.catchDefect((defect) => if (Exit.isSuccess(exit)) return exit.value
defect instanceof TransactionFailure ? Effect.failCause(defect.effectCause as Cause.Cause<E>) : Effect.die(defect), throw new TransactionFailure(exit.cause)
} finally {
txStack.pop()
}
}, config) as A,
).pipe(
Effect.catchDefect((defect) =>
defect instanceof TransactionFailure ? Effect.failCause(defect.effectCause as Cause.Cause<E>) : Effect.die(defect),
),
), ),
), ),
)
return new Proxy(db, {
get(_target, property) {
if (property === "withTransaction") return withTransaction
if (property === "$client") return db.$client
const value = Reflect.get(current(), property)
return typeof value === "function" ? value.bind(current()) : value
},
}) as EffectSQLiteDatabase<TSchema, TRelations> }) as EffectSQLiteDatabase<TSchema, TRelations>
} }

View file

@ -104,22 +104,18 @@ describe("effect drizzle sqlite", () => {
testEffect("runs synchronous Effect programs inside transactions", () => testEffect("runs synchronous Effect programs inside transactions", () =>
Effect.gen(function* () { Effect.gen(function* () {
yield* db.withTransaction((tx) => yield* Effect.gen(function* () {
Effect.gen(function* () { yield* db.insert(users).values({ id: 1, name: "Ada" })
yield* tx.insert(users).values({ id: 1, name: "Ada" }) return yield* db.select().from(users)
return yield* tx.select().from(users) }).pipe(db.withTransaction)
}),
)
expect(yield* db.select().from(users)).toEqual([{ id: 1, name: "Ada" }]) expect(yield* db.select().from(users)).toEqual([{ id: 1, name: "Ada" }])
const exit = yield* Effect.exit( const exit = yield* Effect.exit(
db.withTransaction((tx) => Effect.gen(function* () {
Effect.gen(function* () { yield* db.insert(users).values({ id: 2, name: "Grace" })
yield* tx.insert(users).values({ id: 2, name: "Grace" }) return yield* Effect.fail("rollback")
return yield* Effect.fail("rollback") }).pipe(db.withTransaction),
}),
),
) )
expect(Exit.isFailure(exit)).toBe(true) expect(Exit.isFailure(exit)).toBe(true)
@ -127,6 +123,24 @@ describe("effect drizzle sqlite", () => {
}), }),
) )
testEffect("supports pipeable transactions using the same database service", () =>
Effect.gen(function* () {
const exit = yield* Effect.gen(function* () {
yield* db.insert(users).values({ id: 1, name: "Ada" })
return yield* Effect.fail("rollback")
}).pipe(db.withTransaction, Effect.exit)
expect(Exit.isFailure(exit)).toBe(true)
expect(yield* db.select().from(users)).toEqual([])
yield* Effect.gen(function* () {
yield* db.insert(users).values({ id: 2, name: "Grace" })
}).pipe(db.withTransaction)
expect(yield* db.select().from(users)).toEqual([{ id: 2, name: "Grace" }])
}),
)
testEffect("wraps query failures with query text and parameters", () => testEffect("wraps query failures with query text and parameters", () =>
Effect.gen(function* () { Effect.gen(function* () {
const exit = yield* Effect.exit(db.insert(posts).values({ id: 1, user_id: 404, title: "Missing" })) const exit = yield* Effect.exit(db.insert(posts).values({ id: 1, user_id: 404, title: "Missing" }))