diff --git a/.changeset/fix-atom-batch-listener-writes.md b/.changeset/fix-atom-batch-listener-writes.md new file mode 100644 index 00000000000..f4330ed6049 --- /dev/null +++ b/.changeset/fix-atom-batch-listener-writes.md @@ -0,0 +1,5 @@ +--- +"effect": patch +--- + +Process atom writes queued by batch commit listeners instead of dropping them. diff --git a/packages/effect/src/reactivity/AtomRegistry.ts b/packages/effect/src/reactivity/AtomRegistry.ts index 7269cb96f72..dcda89f6dcb 100644 --- a/packages/effect/src/reactivity/AtomRegistry.ts +++ b/packages/effect/src/reactivity/AtomRegistry.ts @@ -784,11 +784,11 @@ class NodeImpl { } notify(): void { - this.listeners.forEach(notifyListener) - if (batchState.phase === BatchPhase.commit) { batchState.notify.delete(this) } + + this.listeners.forEach(notifyListener) } disposeLifetime(): void { @@ -1113,25 +1113,30 @@ export const batchState = { * @internal */ export function batch(f: () => void): void { + const previousPhase = batchState.phase batchState.phase = BatchPhase.collect batchState.depth++ try { f() if (batchState.depth === 1) { - for (let i = 0; i < batchState.stale.length; i++) { - batchRebuildNode(batchState.stale[i]) - } - batchState.phase = BatchPhase.commit - for (const node of batchState.notify) { - node.notify() - } - batchState.notify.clear() + let i = 0 + do { + batchState.phase = BatchPhase.collect + for (; i < batchState.stale.length; i++) { + batchRebuildNode(batchState.stale[i]) + } + batchState.phase = BatchPhase.commit + for (const node of batchState.notify) { + node.notify() + } + } while (i < batchState.stale.length) } } finally { batchState.depth-- + batchState.phase = previousPhase if (batchState.depth === 0) { - batchState.phase = BatchPhase.disabled batchState.stale = [] + batchState.notify.clear() } } } diff --git a/packages/effect/test/reactivity/Atom.test.ts b/packages/effect/test/reactivity/Atom.test.ts index 4502f19775f..a4f9cf87158 100644 --- a/packages/effect/test/reactivity/Atom.test.ts +++ b/packages/effect/test/reactivity/Atom.test.ts @@ -1083,6 +1083,23 @@ describe("Atom", { concurrent: false }, () => { expect(r.get(derived)).toEqual("2b") }) + it("runs Atom.fn writes from batch commit listeners", () => { + const registry = AtomRegistry.make() + const source = Atom.make(0) + const write = Atom.fn((value: number, get) => Effect.sync(() => get.registry.set(source, value))) + const seen: Array = [] + registry.mount(write) + registry.subscribe(source, (value) => seen.push(value)) + registry.subscribe(source, (value) => { + if (value < 3) registry.set(write, value + 1) + }) + + Atom.batch(() => registry.set(source, 1)) + + assert.deepStrictEqual(seen, [1, 2, 3]) + registry.dispose() + }) + it("initialValues", async () => { const state = Atom.make(0) const r = AtomRegistry.make({