| 659 | * @since 2.0.0 |
| 660 | */ |
| 661 | export const take = <A, E>(self: TxDequeue<A, E>): Effect.Effect<A, E> => |
| 662 | Effect.gen(function*() { |
| 663 | const state = yield* TxRef.get(self.stateRef) |
| 664 | |
| 665 | // Check if queue is done - forward the cause directly |
| 666 | if (state._tag === "Done") { |
| 667 | return yield* Effect.failCause(state.cause) |
| 668 | } |
| 669 | |
| 670 | // If no items available, retry transaction |
| 671 | if (yield* isEmpty(self)) { |
| 672 | return yield* Effect.txRetry |
| 673 | } |
| 674 | |
| 675 | // Take item from queue |
| 676 | const chunk = yield* TxChunk.get(self.items) |
| 677 | const head = Chunk.head(chunk) |
| 678 | if (Option.isNone(head)) { |
| 679 | return yield* Effect.txRetry |
| 680 | } |
| 681 | |
| 682 | yield* TxChunk.drop(self.items, 1) |
| 683 | |
| 684 | // Check if we need to transition Closing → Done |
| 685 | if (state._tag === "Closing" && (yield* isEmpty(self))) { |
| 686 | yield* TxRef.set(self.stateRef, { _tag: "Done", cause: state.cause }) |
| 687 | } |
| 688 | |
| 689 | return head.value |
| 690 | }).pipe(Effect.tx) |
| 691 | |
| 692 | /** |
| 693 | * Tries to take an item from the queue without blocking. |