diff --git a/.changeset/fix-txqueue-offer-all-iterables.md b/.changeset/fix-txqueue-offer-all-iterables.md new file mode 100644 index 00000000000..e6f159cf193 --- /dev/null +++ b/.changeset/fix-txqueue-offer-all-iterables.md @@ -0,0 +1,5 @@ +--- +"effect": patch +--- + +Fix `TxQueue.offerAll` to preserve one-shot iterables across transaction retries and repeated runs. diff --git a/packages/effect/src/TxQueue.ts b/packages/effect/src/TxQueue.ts index 7272128d455..cb611be2aca 100644 --- a/packages/effect/src/TxQueue.ts +++ b/packages/effect/src/TxQueue.ts @@ -608,11 +608,13 @@ export const offerAll: { (self: TxEnqueue, values: Iterable): Effect.Effect> } = dual( 2, - (self: TxEnqueue, values: Iterable): Effect.Effect> => - Effect.gen(function*() { + (self: TxEnqueue, values: Iterable): Effect.Effect> => { + const valuesArray = Array.from(values) + + return Effect.gen(function*() { const rejected: Array = [] - for (const value of values) { + for (const value of valuesArray) { const accepted = yield* offer(self, value) if (!accepted) { rejected.push(value) @@ -621,6 +623,7 @@ export const offerAll: { return rejected }).pipe(Effect.tx) + } ) /** diff --git a/packages/effect/test/TxQueue.test.ts b/packages/effect/test/TxQueue.test.ts index 8c9100ba92c..c041393554d 100644 --- a/packages/effect/test/TxQueue.test.ts +++ b/packages/effect/test/TxQueue.test.ts @@ -230,6 +230,25 @@ describe("TxQueue", () => { assert.strictEqual(size, 5) }))) + it.effect("offerAll preserves a one-shot iterable across retries and runs", () => + Effect.gen(function*() { + const queue = yield* TxQueue.bounded(1) + yield* TxQueue.offer(queue, 1) + + const values = (function*() { + yield 2 + })() + const offer = TxQueue.offerAll(queue, values) + const fiber = yield* Effect.forkChild(offer, { startImmediately: true }) + + assert.strictEqual(yield* TxQueue.take(queue), 1) + assert.deepStrictEqual(yield* Fiber.join(fiber), []) + assert.deepStrictEqual(yield* TxQueue.poll(queue), Option.some(2)) + + assert.deepStrictEqual(yield* offer, []) + assert.deepStrictEqual(yield* TxQueue.poll(queue), Option.some(2)) + })) + it.effect("takeAll works correctly with new signature", () => Effect.tx(Effect.gen(function*() { const queue = yield* TxQueue.bounded(10)