From 41e300ea7198aa0fa6ac61542c23042a17411bdb Mon Sep 17 00:00:00 2001 From: Sebastian Lorenz Date: Sat, 1 Aug 2026 18:27:33 +0000 Subject: [PATCH 1/2] Add reproduction for TxPubSub issue --- packages/effect/test/TxPubSub.test.ts | 28 ++++++++++++++++++++++++++- 1 file changed, 27 insertions(+), 1 deletion(-) diff --git a/packages/effect/test/TxPubSub.test.ts b/packages/effect/test/TxPubSub.test.ts index fe8e1d39060..efe57be4906 100644 --- a/packages/effect/test/TxPubSub.test.ts +++ b/packages/effect/test/TxPubSub.test.ts @@ -1,5 +1,5 @@ import { assert, describe, it } from "@effect/vitest" -import { Effect, Fiber, TxPubSub, TxQueue } from "effect" +import { Effect, Fiber, Option, TxPubSub, TxQueue } from "effect" describe("TxPubSub", () => { describe("constructors", () => { @@ -95,6 +95,32 @@ describe("TxPubSub", () => { ) })) + it.effect("publishAll preserves a one-shot iterable across retries", () => + Effect.gen(function*() { + const hub = yield* TxPubSub.bounded(1) + + yield* Effect.scoped( + Effect.gen(function*() { + const sub = yield* TxPubSub.subscribe(hub) + yield* TxPubSub.publish(hub, 1) + + let iterations = 0 + const values = (function*() { + iterations++ + yield 2 + })() + const fiber = yield* Effect.forkChild(TxPubSub.publishAll(hub, values)) + while (iterations === 0) { + yield* Effect.yieldNow + } + + assert.strictEqual(yield* TxQueue.take(sub), 1) + assert.strictEqual(yield* Fiber.join(fiber), true) + assert.deepStrictEqual(yield* TxQueue.poll(sub), Option.some(2)) + }) + ) + })) + it.effect("subscriber only receives messages published after subscription", () => Effect.gen(function*() { const hub = yield* Effect.tx(TxPubSub.unbounded()) From 8c32aef0c5320411d1d5e0ada19a818b75d5e2be Mon Sep 17 00:00:00 2001 From: Tim Smart Date: Mon, 3 Aug 2026 10:40:17 +1200 Subject: [PATCH 2/2] Fix TxPubSub publishAll iterable retries --- .changeset/fix-txpubsub-publish-all-iterables.md | 5 +++++ packages/effect/src/TxPubSub.ts | 9 ++++++--- packages/effect/test/TxPubSub.test.ts | 3 ++- 3 files changed, 13 insertions(+), 4 deletions(-) create mode 100644 .changeset/fix-txpubsub-publish-all-iterables.md diff --git a/.changeset/fix-txpubsub-publish-all-iterables.md b/.changeset/fix-txpubsub-publish-all-iterables.md new file mode 100644 index 00000000000..390063b2ba9 --- /dev/null +++ b/.changeset/fix-txpubsub-publish-all-iterables.md @@ -0,0 +1,5 @@ +--- +"effect": patch +--- + +Fix `TxPubSub.publishAll` dropping values from one-shot iterables when a transaction retries. diff --git a/packages/effect/src/TxPubSub.ts b/packages/effect/src/TxPubSub.ts index 165fb73f454..7f8b327104a 100644 --- a/packages/effect/src/TxPubSub.ts +++ b/packages/effect/src/TxPubSub.ts @@ -9,6 +9,7 @@ * * @since 4.0.0 */ +import * as Arr from "./Array.ts" import * as Effect from "./Effect.ts" import { dual } from "./Function.ts" import type { Inspectable } from "./Inspectable.ts" @@ -471,17 +472,19 @@ export const publishAll: { (self: TxPubSub, values: Iterable): Effect.Effect } = dual( 2, - (self: TxPubSub, values: Iterable): Effect.Effect => - Effect.gen(function*() { + (self: TxPubSub, values: Iterable): Effect.Effect => { + const valuesArray = Arr.fromIterable(values) + return Effect.gen(function*() { if (yield* TxRef.get(self.shutdownRef)) return false let allAccepted = true - for (const value of values) { + for (const value of valuesArray) { const accepted = yield* publish(self, value) if (!accepted) allAccepted = false } return allAccepted }).pipe(Effect.tx) + } ) /** diff --git a/packages/effect/test/TxPubSub.test.ts b/packages/effect/test/TxPubSub.test.ts index efe57be4906..20763d432c9 100644 --- a/packages/effect/test/TxPubSub.test.ts +++ b/packages/effect/test/TxPubSub.test.ts @@ -105,12 +105,13 @@ describe("TxPubSub", () => { yield* TxPubSub.publish(hub, 1) let iterations = 0 + const hasIterated = () => iterations > 0 const values = (function*() { iterations++ yield 2 })() const fiber = yield* Effect.forkChild(TxPubSub.publishAll(hub, values)) - while (iterations === 0) { + while (!hasIterated()) { yield* Effect.yieldNow }