Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions .changeset/fix-txqueue-closing-drain.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
"effect": patch
---

Ensure `TxQueue.poll` and `TxQueue.clear` complete a closing queue after draining its buffered items.
13 changes: 10 additions & 3 deletions packages/effect/src/TxQueue.ts
Original file line number Diff line number Diff line change
Expand Up @@ -725,6 +725,11 @@ export const poll = <A, E>(self: TxDequeue<A, E>): Effect.Effect<Option.Option<A
}

yield* TxChunk.drop(self.items, 1)

if (state._tag === "Closing" && (yield* isEmpty(self))) {
yield* TxRef.set(self.stateRef, { _tag: "Done", cause: state.cause })
}

return Option.some(head.value)
}).pipe(Effect.tx)

Expand Down Expand Up @@ -1269,12 +1274,11 @@ export const end = <A, E>(self: TxEnqueue<A, E | Cause.Done>): Effect.Effect<boo
failCause(self, Cause.fail(Cause.Done()))

/**
* Removes and returns all currently buffered elements without changing the
* queue state.
* Removes and returns all currently buffered elements.
*
* **Details**
*
* If the queue is already done with a `Cause.Done` error, returns an empty array. If the queue is done for any other cause, including interruption or failure, that cause is propagated.
* If the queue is closing, draining its buffered elements transitions it to done. If the queue is already done with a `Cause.Done` error, returns an empty array. If the queue is done for any other cause, including interruption or failure, that cause is propagated.
*
* **Example** (Clearing queues)
*
Expand Down Expand Up @@ -1311,6 +1315,9 @@ export const clear = <A, E>(self: TxEnqueue<A, E>): Effect.Effect<Array<A>, 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)

Expand Down
20 changes: 20 additions & 0 deletions packages/effect/test/TxQueue.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<number>(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<number>(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<number>(10)
Expand Down
Loading