From ce7d5f0914c77fcf100799c8c207f456f64f0234 Mon Sep 17 00:00:00 2001 From: Sebastian Lorenz Date: Sat, 1 Aug 2026 18:36:21 +0000 Subject: [PATCH 1/3] Add reproduction for TxQueue issue --- packages/effect/test/TxQueue.test.ts | 20 ++++++++++++++++++++ 1 file changed, 20 insertions(+) diff --git a/packages/effect/test/TxQueue.test.ts b/packages/effect/test/TxQueue.test.ts index 8c9100ba92c..f48bea8690e 100644 --- a/packages/effect/test/TxQueue.test.ts +++ b/packages/effect/test/TxQueue.test.ts @@ -230,6 +230,26 @@ describe("TxQueue", () => { assert.strictEqual(size, 5) }))) + it.effect("offerAll preserves a one-shot iterable across retries", () => + Effect.gen(function*() { + const queue = yield* TxQueue.bounded(1) + yield* TxQueue.offer(queue, 1) + + let iterations = 0 + const values = (function*() { + iterations++ + yield 2 + })() + const fiber = yield* Effect.forkChild(TxQueue.offerAll(queue, values)) + while (iterations === 0) { + yield* Effect.yieldNow + } + + assert.strictEqual(yield* TxQueue.take(queue), 1) + assert.deepStrictEqual(yield* Fiber.join(fiber), []) + 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) From 2afea5a1021244caa0a3793aba1d0c444f9bcbee Mon Sep 17 00:00:00 2001 From: Tim Smart Date: Mon, 3 Aug 2026 10:40:16 +1200 Subject: [PATCH 2/3] Fix TxQueue offerAll iterable retries --- packages/effect/src/TxQueue.ts | 22 +++++++++++++--------- packages/effect/test/TxQueue.test.ts | 10 ++++------ 2 files changed, 17 insertions(+), 15 deletions(-) diff --git a/packages/effect/src/TxQueue.ts b/packages/effect/src/TxQueue.ts index 7272128d455..5ac537fc600 100644 --- a/packages/effect/src/TxQueue.ts +++ b/packages/effect/src/TxQueue.ts @@ -609,18 +609,22 @@ export const offerAll: { } = dual( 2, (self: TxEnqueue, values: Iterable): Effect.Effect> => - Effect.gen(function*() { - const rejected: Array = [] + Effect.suspend(() => { + const valuesArray = Array.from(values) + + return Effect.gen(function*() { + const rejected: Array = [] - for (const value of values) { - const accepted = yield* offer(self, value) - if (!accepted) { - rejected.push(value) + for (const value of valuesArray) { + const accepted = yield* offer(self, value) + if (!accepted) { + rejected.push(value) + } } - } - return rejected - }).pipe(Effect.tx) + return rejected + }).pipe(Effect.tx) + }) ) /** diff --git a/packages/effect/test/TxQueue.test.ts b/packages/effect/test/TxQueue.test.ts index f48bea8690e..926bf6fe51a 100644 --- a/packages/effect/test/TxQueue.test.ts +++ b/packages/effect/test/TxQueue.test.ts @@ -1,5 +1,5 @@ import { assert, describe, it } from "@effect/vitest" -import { Cause, Effect, Fiber, Option, Result, TxQueue } from "effect" +import { Cause, Effect, Fiber, Latch, Option, Result, TxQueue } from "effect" describe("TxQueue", () => { describe("interfaces", () => { @@ -235,15 +235,13 @@ describe("TxQueue", () => { const queue = yield* TxQueue.bounded(1) yield* TxQueue.offer(queue, 1) - let iterations = 0 + const iterationStarted = Latch.makeUnsafe() const values = (function*() { - iterations++ + iterationStarted.openUnsafe() yield 2 })() const fiber = yield* Effect.forkChild(TxQueue.offerAll(queue, values)) - while (iterations === 0) { - yield* Effect.yieldNow - } + yield* iterationStarted.await assert.strictEqual(yield* TxQueue.take(queue), 1) assert.deepStrictEqual(yield* Fiber.join(fiber), []) From a8fad6e83a30036cbbeccdef42eda36f51c027ff Mon Sep 17 00:00:00 2001 From: Tim Smart Date: Mon, 3 Aug 2026 10:51:45 +1200 Subject: [PATCH 3/3] Address TxQueue offerAll review feedback --- .changeset/fix-txqueue-offer-all-iterables.md | 5 ++++ packages/effect/src/TxQueue.ts | 25 +++++++++---------- packages/effect/test/TxQueue.test.ts | 13 +++++----- 3 files changed, 24 insertions(+), 19 deletions(-) create mode 100644 .changeset/fix-txqueue-offer-all-iterables.md 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 5ac537fc600..cb611be2aca 100644 --- a/packages/effect/src/TxQueue.ts +++ b/packages/effect/src/TxQueue.ts @@ -608,23 +608,22 @@ export const offerAll: { (self: TxEnqueue, values: Iterable): Effect.Effect> } = dual( 2, - (self: TxEnqueue, values: Iterable): Effect.Effect> => - Effect.suspend(() => { - const valuesArray = Array.from(values) + (self: TxEnqueue, values: Iterable): Effect.Effect> => { + const valuesArray = Array.from(values) - return Effect.gen(function*() { - const rejected: Array = [] + return Effect.gen(function*() { + const rejected: Array = [] - for (const value of valuesArray) { - const accepted = yield* offer(self, value) - if (!accepted) { - rejected.push(value) - } + for (const value of valuesArray) { + const accepted = yield* offer(self, value) + if (!accepted) { + rejected.push(value) } + } - return rejected - }).pipe(Effect.tx) - }) + return rejected + }).pipe(Effect.tx) + } ) /** diff --git a/packages/effect/test/TxQueue.test.ts b/packages/effect/test/TxQueue.test.ts index 926bf6fe51a..c041393554d 100644 --- a/packages/effect/test/TxQueue.test.ts +++ b/packages/effect/test/TxQueue.test.ts @@ -1,5 +1,5 @@ import { assert, describe, it } from "@effect/vitest" -import { Cause, Effect, Fiber, Latch, Option, Result, TxQueue } from "effect" +import { Cause, Effect, Fiber, Option, Result, TxQueue } from "effect" describe("TxQueue", () => { describe("interfaces", () => { @@ -230,22 +230,23 @@ describe("TxQueue", () => { assert.strictEqual(size, 5) }))) - it.effect("offerAll preserves a one-shot iterable across retries", () => + 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 iterationStarted = Latch.makeUnsafe() const values = (function*() { - iterationStarted.openUnsafe() yield 2 })() - const fiber = yield* Effect.forkChild(TxQueue.offerAll(queue, values)) - yield* iterationStarted.await + 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", () =>