From 0719851887bc20bb9e3c92778a51bea6932c682b Mon Sep 17 00:00:00 2001 From: Sebastian Lorenz Date: Sat, 1 Aug 2026 18:31:31 +0000 Subject: [PATCH 1/2] 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..4f8905afd02 100644 --- a/packages/effect/test/TxQueue.test.ts +++ b/packages/effect/test/TxQueue.test.ts @@ -219,6 +219,26 @@ describe("TxQueue", () => { assert.deepStrictEqual(maybe, Option.some(42)) }))) + it.effect("poll completes a closing queue after removing the last item", () => + Effect.tx(Effect.gen(function*() { + const queue = yield* TxQueue.bounded(1) + yield* TxQueue.offer(queue, 42) + yield* TxQueue.interrupt(queue) + + assert.deepStrictEqual(yield* TxQueue.poll(queue), Option.some(42)) + assert.strictEqual(yield* TxQueue.isDone(queue), true) + }))) + + it.effect("clear completes a closing queue after removing all items", () => + Effect.tx(Effect.gen(function*() { + const queue = yield* TxQueue.bounded(1) + yield* TxQueue.offer(queue, 42) + yield* TxQueue.interrupt(queue) + + assert.deepStrictEqual(yield* TxQueue.clear(queue), [42]) + assert.strictEqual(yield* TxQueue.isDone(queue), true) + }))) + it.effect("offerAll works correctly", () => Effect.tx(Effect.gen(function*() { const queue = yield* TxQueue.bounded(10) From 9f6eac19ea64965e90c11ea2f654d59bfec0463c Mon Sep 17 00:00:00 2001 From: Tim Smart Date: Mon, 3 Aug 2026 11:21:41 +1200 Subject: [PATCH 2/2] Fix TxQueue closing drains --- .changeset/fix-txqueue-closing-drain.md | 5 +++++ packages/effect/src/TxQueue.ts | 13 ++++++++++--- 2 files changed, 15 insertions(+), 3 deletions(-) create mode 100644 .changeset/fix-txqueue-closing-drain.md diff --git a/.changeset/fix-txqueue-closing-drain.md b/.changeset/fix-txqueue-closing-drain.md new file mode 100644 index 00000000000..41b0208f88c --- /dev/null +++ b/.changeset/fix-txqueue-closing-drain.md @@ -0,0 +1,5 @@ +--- +"effect": patch +--- + +Ensure `TxQueue.poll` and `TxQueue.clear` complete a closing queue after draining its buffered items. diff --git a/packages/effect/src/TxQueue.ts b/packages/effect/src/TxQueue.ts index 7272128d455..37a44addf61 100644 --- a/packages/effect/src/TxQueue.ts +++ b/packages/effect/src/TxQueue.ts @@ -725,6 +725,11 @@ export const poll = (self: TxDequeue): Effect.Effect(self: TxEnqueue): Effect.Effect(self: TxEnqueue): Effect.Effect, Excl } const chunk = yield* TxChunk.get(self.items) yield* TxChunk.set(self.items, Chunk.empty()) + if (state._tag === "Closing") { + yield* TxRef.set(self.stateRef, { _tag: "Done", cause: state.cause }) + } return Chunk.toArray(chunk) }).pipe(Effect.tx)