Queue
Passes values asynchronously between fibers.
A Queue<A, E> accepts values, hands each value to one consumer in offer
order, and can complete, fail, interrupt, or shut down. Queues can be bounded
or unbounded, and bounded queues can suspend, drop, or slide values when
producers are faster than consumers.
Completion
Signals queue completion.
When to use
Use to stop accepting new offers while allowing already queued messages to be consumed.
Details
Returns false if the queue is already done.
Signature
declare function end<A, E>(self: Enqueue<A, Done<void> | E>): Effect<boolean>Example
(Ending queues)
import { Cause, Effect, Queue } from "effect"
const program = Effect.gen(function*() { const queue = yield* Queue.bounded<number, Cause.Done>(10)
// Add some messages yield* Queue.offer(queue, 1) yield* Queue.offer(queue, 2)
// Signal completion - no more messages will be accepted const ended = yield* Queue.end(queue)
// Trying to offer more messages will return false const offerResult = yield* Queue.offer(queue, 3)
// But we can still take existing messages const message = yield* Queue.take(queue) return [ended, offerResult, message]})
await Effect.runPromise(program) // => [true, false, 1]Signals queue completion synchronously.
When to use
Use when implementing low-level queue integrations that must complete a queue
without wrapping the operation in Effect.
Details
Returns false if the queue is already done.
Gotchas
This is an unsafe operation that directly modifies the queue without Effect wrapping.
Signature
declare function endUnsafe<A, E>(self: Enqueue<A, Done<void> | E>): booleanExample
(Ending queues synchronously)
import { Cause, Effect, Queue } from "effect"
// Create a queue and use unsafe operationsconst program = Effect.gen(function*() { const queue = yield* Queue.bounded<number, Cause.Done>(10)
// Add some messages Queue.offerUnsafe(queue, 1) Queue.offerUnsafe(queue, 2)
// End the queue synchronously const ended = Queue.endUnsafe(queue)
// Existing messages can still be consumed while the queue is closing const states = [queue.state._tag]
Queue.takeUnsafe(queue) Queue.takeUnsafe(queue)
// After buffered messages are consumed, the queue is done states.push(queue.state._tag) return { ended, states }})
await Effect.runPromise(program) // => { ended: true, states: ["Closing", "Done"] }Fails the queue with an error. If the queue is already done, false is
returned.
Signature
declare function fail<A, E>(self: Enqueue<A, E>, error: E): Effect<boolean, never, never>Example
(Failing queues with an error)
import { Effect, Exit, Queue } from "effect"
const program = Effect.gen(function*() { const queue = yield* Queue.bounded<number, string>(10)
// Fail the queue with an error const failed = yield* Queue.fail(queue, "Something went wrong")
// Taking from the failed queue fails with the error const exit = yield* Effect.exit(Queue.take(queue)) return [failed, exit]})
await Effect.runPromise(program) // => [true, Exit.fail("Something went wrong")]Fails the queue with a cause. If the queue is already done, false is
returned.
Signature
declare const failCause: { <E>(cause: Cause<E>): <A>(self: Enqueue<A, E>) => Effect<boolean>; <A, E>(self: Enqueue<A, E>, cause: Cause<E>): Effect<boolean>;}Example
(Failing queues with a cause)
import { Cause, Effect, Exit, Queue } from "effect"
const program = Effect.gen(function*() { const queue = yield* Queue.bounded<number, string>(10)
// Create a cause and fail the queue const cause = Cause.fail("Queue processing failed") const failed = yield* Queue.failCause(queue, cause)
// The queue is now done with the specified failure cause const exit = yield* Effect.exit(Queue.take(queue)) return [failed, exit]})
await Effect.runPromise(program) // => [true, Exit.failCause(Cause.fail("Queue processing failed"))]failCauseUnsafe
Fails the queue with a cause synchronously. If the queue is already done, false is
returned.
When to use
Use when queue completion must be driven from synchronous internals while
preserving the full failure Cause.
Gotchas
This is an unsafe operation that directly modifies the queue without Effect wrapping.
Signature
declare function failCauseUnsafe<A, E>(self: Enqueue<A, E>, cause: Cause<E>): booleanExample
(Failing queues with a cause synchronously)
import { Cause, Effect, Exit, Queue } from "effect"
const program = Effect.gen(function*() { const queue = yield* Queue.bounded<number, string>(10)
// Create a cause and fail the queue synchronously const cause = Cause.fail("Processing error") const failed = Queue.failCauseUnsafe(queue, cause)
// The queue is now done with the specified failure cause const exit = Queue.takeUnsafe(queue) return [failed, exit]})
await Effect.runPromise(program) // => [true, Exit.failCause(Cause.fail("Processing error"))]Interrupts the queue gracefully, transitioning it to a closing state.
Details
This operation stops accepting new offers but allows existing messages to be consumed. Once all messages are drained, the queue transitions to the Done state with an interrupt cause.
Signature
declare function interrupt<A, E>(self: Enqueue<A, E>): Effect<boolean>Example
(Interrupting queues gracefully)
import { Cause, Effect, Queue } from "effect"
const program = Effect.gen(function*() { const queue = yield* Queue.bounded<number>(10)
// Add some messages yield* Queue.offer(queue, 1) yield* Queue.offer(queue, 2)
// Interrupt gracefully - no more offers accepted, but messages can be consumed const interrupted = yield* Queue.interrupt(queue)
// Trying to offer more messages will return false const offerResult = yield* Queue.offer(queue, 3)
// But we can still take existing messages const message1 = yield* Queue.take(queue)
const message2 = yield* Queue.take(queue)
// After all messages are consumed, queue is done const isDone = queue.state._tag === "Done" return { interrupted, offerResult, messages: [message1, message2], isDone }})
await Effect.runPromise(program) // => { interrupted: true, offerResult: false, messages: [1, 2], isDone: true }Runs an Effect into a Queue, where success ends the queue and failure
fails the queue.
Signature
declare const into: { <A, E>(self: Enqueue<A, Done<void> | E>): <AX, EX, RX>(effect: Effect<AX, EX, RX>) => Effect<boolean, never, RX>; <AX, E, EX, RX, A>(effect: Effect<AX, EX, RX>, self: Enqueue<A, Done<void> | E>): Effect<boolean, never, RX>;}Example
(Running effects into queues)
import { Cause, Effect, Exit, Queue } from "effect"
const program = Effect.gen(function*() { const queue = yield* Queue.bounded<number, Cause.Done>(10)
// Create an effect that succeeds const dataProcessing = Effect.gen(function*() { yield* Effect.yieldNow return "Processing completed successfully" })
// Pipe the effect into the queue // If dataProcessing succeeds, queue ends successfully // If dataProcessing fails, queue fails with the error const effectIntoQueue = Queue.into(queue)(dataProcessing)
const wasCompleted = yield* effectIntoQueue const exit = yield* Effect.exit(Queue.take(queue)) return [wasCompleted, exit]})
await Effect.runPromise(program) // => [true, Exit.fail(Cause.Done())]Shuts down the queue immediately, discarding buffered messages and resuming pending operations.
Details
The operation is idempotent and returns true, including when the queue has
already been shut down or completed.
Signature
declare function shutdown<A, E>(self: Enqueue<A, E>): Effect<boolean>Example
(Shutting down queues)
import { Effect, Queue } from "effect"
const program = Effect.gen(function*() { const queue = yield* Queue.bounded<number>(2)
// Add messages yield* Queue.offer(queue, 1) yield* Queue.offer(queue, 2)
// Shutdown clears buffered messages and prevents further offers const wasShutdown = yield* Queue.shutdown(queue)
// Queue is now done and cleared const size = yield* Queue.size(queue) return { wasShutdown, size }})
await Effect.runPromise(program) // => { wasShutdown: true, size: 0 }Constructors
Creates a bounded queue with the specified capacity that uses backpressure strategy.
Details
When the queue reaches capacity, producers will be suspended until space becomes available. This ensures all messages are processed but may slow down producers.
Signature
declare function bounded<A, E = never>(capacity: number): Effect<Queue<A, E>>Example
(Creating bounded queues)
import { Effect, Queue } from "effect"
const program = Effect.gen(function*() { const queue = yield* Queue.bounded<string>(5)
// This will succeed as queue has capacity yield* Queue.offer(queue, "first") yield* Queue.offer(queue, "second")
const size = yield* Queue.size(queue) return size})
await Effect.runPromise(program) // => 2Creates a bounded queue with dropping strategy. When the queue reaches capacity, new elements are dropped and the offer operation returns false.
When to use
Use when you need producer offers not to block while preserving existing queued messages, even if new messages may be dropped when the queue is full.
Signature
declare function dropping<A, E = never>(capacity: number): Effect<Queue<A, E>>Example
(Creating dropping queues)
import { Effect, Queue } from "effect"
const program = Effect.gen(function*() { const queue = yield* Queue.dropping<number>(2)
// Fill the queue to capacity const success1 = yield* Queue.offer(queue, 1) const success2 = yield* Queue.offer(queue, 2)
// This will be dropped const success3 = yield* Queue.offer(queue, 3)
const all = yield* Queue.takeAll(queue) return [success1, success2, success3, all]})
await Effect.runPromise(program) // => [true, true, false, [1, 2]]Creates a Queue with optional capacity and overflow strategy.
Details
By default the queue is unbounded and uses the "suspend" strategy. Provide
capacity for a bounded queue and choose "suspend", "dropping", or
"sliding" to control what happens when the queue is full. The returned
queue can be offered to, taken from, failed, ended, interrupted, or shut down.
Signature
declare function make<A, E = never>(options?: { readonly capacity?: number; readonly strategy?: "sliding" | "dropping" | "suspend";}): Effect<Queue<A, E>>Example
(Creating queues)
import { Cause, Effect, Queue } from "effect"
const program = Effect.gen(function*() { const queue = yield* Queue.make<number, string | Cause.Done>()
// add messages to the queue yield* Queue.offer(queue, 1) yield* Queue.offer(queue, 2) yield* Queue.offerAll(queue, [3, 4, 5])
// take messages from the queue const messages = yield* Queue.takeAll(queue)
// signal that the queue is done yield* Queue.end(queue) const done = yield* Effect.flip(Queue.take(queue))
// signal that another queue has failed const failedQueue = yield* Queue.make<number, string>() const failed = yield* Queue.fail(failedQueue, "boom") return { messages, done, failed }})
await Effect.runPromise(program) // => { messages: [1, 2, 3, 4, 5], done: Cause.Done(), failed: true }Creates a bounded queue with sliding strategy. When the queue reaches capacity, new elements are added and the oldest elements are dropped.
When to use
Use when you need producer offers not to block and can accept dropping the oldest messages, such as when maintaining a rolling window of recent values.
Signature
declare function sliding<A, E = never>(capacity: number): Effect<Queue<A, E>>Example
(Creating sliding queues)
import { Effect, Queue } from "effect"
const program = Effect.gen(function*() { const queue = yield* Queue.sliding<number>(3)
// Fill the queue to capacity yield* Queue.offer(queue, 1) yield* Queue.offer(queue, 2) yield* Queue.offer(queue, 3)
// This will succeed, dropping the oldest element (1) yield* Queue.offer(queue, 4)
const all = yield* Queue.takeAll(queue) return all})
await Effect.runPromise(program) // => [2, 3, 4]Creates an unbounded queue that can grow to any size without blocking producers.
When to use
Use when you need producers to add messages without backpressure and accept unbounded memory growth.
Signature
declare function unbounded<A, E = never>(): Effect<Queue<A, E>>Example
(Creating unbounded queues)
import { Effect, Queue } from "effect"
const program = Effect.gen(function*() { const queue = yield* Queue.unbounded<string>()
// Producers can always add messages without blocking yield* Queue.offer(queue, "message1") yield* Queue.offer(queue, "message2") yield* Queue.offerAll(queue, ["message3", "message4", "message5"])
// Check current size const size = yield* Queue.size(queue)
// Take all messages const messages = yield* Queue.takeAll(queue) return { size, messages }})
await Effect.runPromise(program) // => { size: 5, messages: ["message1", "message2", "message3", "message4", "message5"] }Converting
Narrows a Queue to a Dequeue, exposing the consumer side of the queue.
When to use
Use to pass a queue to code that should consume values while keeping producer-side operations out of that code's TypeScript type.
Gotchas
This is a type-level narrowing operation. It returns the same queue object and does not create a runtime wrapper.
See
Signature
declare const asDequeue: <A, E>(self: Queue<A, E>) => Dequeue<A, E>Converts a Queue to its write-only Enqueue interface.
When to use
Use to expose only the producer side of a Queue to code that should offer
values or signal queue lifecycle.
Gotchas
This is a type-level capability restriction. It returns the same queue object, so it does not hide read operations at runtime.
See
Signature
declare function asEnqueue<A, E>(self: Queue<A, E>): Enqueue<A, E>Guards
Type guard to check if a value is a Dequeue.
When to use
Use to narrow an unknown value before passing it to read-side queue operations.
See
Signature
declare function isDequeue<A = unknown, E = unknown>(u: unknown): u is Dequeue<A, E>Type guard to check if a value is an Enqueue.
When to use
Use to narrow an unknown value before calling queue operations that require write-side access.
Gotchas
A full Queue also satisfies this guard because every queue includes the
enqueue side.
See
Signature
declare function isEnqueue<A = unknown, E = unknown>(u: unknown): u is Enqueue<A, E>Type guard to check if a value is a Queue.
When to use
Use to narrow an unknown value to a full Queue before passing it to APIs
that need both offering and taking capabilities.
See
Signature
declare function isQueue<A = unknown, E = unknown>(u: unknown): u is Queue<A, E>Models
A Dequeue is a queue that can be taken from.
Details
This interface represents the read-only part of a Queue, allowing you to take elements from the queue but not offer elements to it.
Signature
interface Dequeue<out A, out E = never> extends Inspectable { readonly "~effect/Queue/Dequeue": Variance<A, E>; capacity: number; readonly dispatcher: SchedulerDispatcher; messages: MutableList<any>; scheduleRunning: boolean; state: State<any, any>; readonly strategy: "sliding" | "dropping" | "suspend";}Example
(Taking through dequeue handles)
import { Effect, Queue } from "effect"
const program = Effect.gen(function*() { const queue = yield* Queue.bounded<string, never>(10)
// A Dequeue can only take elements const dequeue: Queue.Dequeue<string> = queue
// Pre-populate the queue yield* Queue.offerAll(queue, ["a", "b", "c"])
// Take elements using dequeue interface const item = yield* Queue.take(dequeue) return item})
await Effect.runPromise(program) // => "a"An Enqueue is a queue that can be offered to.
Details
This interface represents the write-only part of a Queue, allowing you to offer elements to the queue but not take elements from it.
Signature
interface Enqueue<in A, in E = never> extends Inspectable { readonly "~effect/Queue/Enqueue": Variance<A, E>; capacity: number; readonly dispatcher: SchedulerDispatcher; messages: MutableList<any>; scheduleRunning: boolean; state: State<any, any>; readonly strategy: "sliding" | "dropping" | "suspend";}Example
(Offering through enqueue handles)
import { Effect, Queue } from "effect"
// Function that only needs write access to a queueconst producer = (enqueue: Queue.Enqueue<string>) => Effect.gen(function*() { yield* Queue.offer(enqueue, "hello") yield* Queue.offerAll(enqueue, ["world", "!"]) })
const program = Effect.gen(function*() { const queue = yield* Queue.bounded<string>(10) yield* producer(queue) return yield* Queue.takeAll(queue)})
await Effect.runPromise(program) // => ["hello", "world", "!"]A Queue is an asynchronous queue that can be offered to and taken from.
Details
It also supports signaling that it is done or failed.
Signature
interface Queue<in out A, in out E = never> extends Enqueue<A, E>, Dequeue<A, E> { readonly "~effect/Queue": Variance<A, E>;}Example
(Offering and taking queue values)
import { Effect, Queue } from "effect"
const program = Effect.gen(function*() { // Create a bounded queue const queue = yield* Queue.bounded<string>(10)
// Producer: offer items to the queue yield* Queue.offer(queue, "hello") yield* Queue.offerAll(queue, ["world", "!"])
// Consumer: take items from the queue const item1 = yield* Queue.take(queue) const item2 = yield* Queue.take(queue) const item3 = yield* Queue.take(queue)
return [item1, item2, item3]})
await Effect.runPromise(program) // => ["hello", "world", "!"]Offering
Adds a message to the queue. Returns false if the queue is done.
Details
For bounded queues, this operation may suspend if the queue is at capacity, depending on the backpressure strategy. For dropping/sliding queues, it may return false or succeed immediately by dropping/sliding existing messages.
Signature
declare function offer<A, E>(self: Enqueue<A, E>, message: NoInfer<A>): Effect<boolean>Example
(Offering a value)
import { Effect, Queue } from "effect"
const program = Effect.gen(function*() { const queue = yield* Queue.bounded<number>(3)
// Successfully add messages to queue const success1 = yield* Queue.offer(queue, 1) const success2 = yield* Queue.offer(queue, 2)
// Queue state const size = yield* Queue.size(queue) return { offered: [success1, success2], size }})
await Effect.runPromise(program) // => { offered: [true, true], size: 2 }Adds multiple messages to the queue. Returns the remaining messages that were not added.
When to use
Use when producers can submit a batch at once and need to know which messages did not fit under the queue's capacity strategy.
Details
For bounded queues, this operation may suspend if the queue doesn't have enough capacity. The operation returns an array of messages that couldn't be added (empty array means all messages were successfully added).
Signature
declare function offerAll<A, E>(self: Enqueue<A, E>, messages: Iterable<A>): Effect<Array<A>>Example
(Offering multiple values)
import { Effect, Queue } from "effect"
const program = Effect.gen(function*() { const queue = yield* Queue.dropping<number>(3)
// Try to add more messages than capacity without suspending const remaining1 = yield* Queue.offerAll(queue, [1, 2, 3, 4, 5]) return remaining1})
await Effect.runPromise(program) // => [4, 5]offerAllUnsafe
Adds multiple messages to the queue synchronously. Returns the remaining messages that were not added.
When to use
Use when queue internals or a performance boundary need a synchronous batch offer and can handle any messages that do not fit.
Gotchas
This is an unsafe operation that directly modifies the queue without Effect wrapping.
Signature
declare function offerAllUnsafe<A, E>(self: Enqueue<A, E>, messages: Iterable<A>): Array<A>Example
(Offering multiple values synchronously)
import { Effect, Queue } from "effect"
// Create a bounded queue and use unsafe APIconst program = Effect.gen(function*() { const queue = yield* Queue.bounded<number>(3)
// Try to add 5 messages to capacity-3 queue using unsafe API const remaining = Queue.offerAllUnsafe(queue, [1, 2, 3, 4, 5])
// Check what's in the queue const size = Queue.sizeUnsafe(queue) return { remaining, size }})
await Effect.runPromise(program) // => { remaining: [4, 5], size: 3 }offerUnsafe
Adds a message to the queue synchronously. Returns false if the queue is done.
When to use
Use when you are already in synchronous queue internals or a performance
boundary where wrapping the mutation in Effect is intentionally avoided.
Gotchas
This is an unsafe operation that directly modifies the queue without Effect wrapping. Use this only when you're certain about the synchronous nature of the operation.
Signature
declare function offerUnsafe<A, E>(self: Enqueue<A, E>, message: NoInfer<A>): booleanExample
(Offering a value synchronously)
import { Effect, Queue } from "effect"
// Create a queue effect and extract the queue for unsafe operationsconst program = Effect.gen(function*() { const queue = yield* Queue.bounded<number>(3)
// Add messages synchronously using unsafe API const success1 = Queue.offerUnsafe(queue, 1) const success2 = Queue.offerUnsafe(queue, 2)
// Check current size const size = Queue.sizeUnsafe(queue) return { offered: [success1, success2], size }})
await Effect.runPromise(program) // => { offered: [true, true], size: 2 }Other
Signature
declare function await<A, E>(self: Dequeue<A, E>): Effect<void, Exclude<E, Done<void>>>Companion namespace containing type-level metadata for the Dequeue
read-only queue interface.
Companion namespace containing type-level metadata for the Enqueue
write-only queue interface.
Companion namespace containing type-level metadata and low-level state types
for Queue.
Predicates
Checks whether the queue is full.
Signature
declare function isFull<A, E>(self: Dequeue<A, E>): Effect<boolean>Example
(Checking if queues are full)
import { Cause, Effect, Queue } from "effect"
const program = Effect.gen(function*() { const queue = yield* Queue.bounded<number, Cause.Done>(3)
const before = yield* Queue.isFull(queue)
// Add some messages yield* Queue.offerAll(queue, [1, 2, 3])
const after = yield* Queue.isFull(queue) return [before, after]})
await Effect.runPromise(program) // => [false, true]isFullUnsafe
Checks whether the queue is full synchronously.
When to use
Use when an immediate Queue capacity snapshot is needed outside effectful
code and racing queue changes are acceptable.
Signature
declare function isFullUnsafe<A, E>(self: Dequeue<A, E>): booleanExample
(Checking fullness synchronously)
import { Cause, Effect, Queue } from "effect"
const program = Effect.gen(function*() { const queue = yield* Queue.bounded<number, Cause.Done>(3)
const before = Queue.isFullUnsafe(queue)
// Add some messages yield* Queue.offerAll(queue, [1, 2, 3])
const after = Queue.isFullUnsafe(queue) return [before, after]})
await Effect.runPromise(program) // => [false, true]Sizes
Returns the current number of buffered messages in the queue.
Details
After end, a queue remains Closing while buffered messages are drained,
and its size continues to include those messages. A Done queue reports a
size of 0.
Signature
declare function size<A, E>(self: Dequeue<A, E>): Effect<number>Example
(Checking queue size)
import { Cause, Effect, Queue } from "effect"
const program = Effect.gen(function*() { const queue = yield* Queue.bounded<number, Cause.Done>(10)
// Check size of empty queue const size1 = yield* Queue.size(queue)
// Add some messages yield* Queue.offerAll(queue, [1, 2, 3, 4, 5])
// Check size after adding messages const size2 = yield* Queue.size(queue)
// End the queue yield* Queue.end(queue)
// Ending retains the buffered size while the queue is Closing const size3 = yield* Queue.size(queue) return [size1, size2, size3]})
await Effect.runPromise(program) // => [0, 5, 5]sizeUnsafe
Returns the current number of buffered messages in the queue synchronously.
When to use
Use when you need an immediate Queue size snapshot for diagnostics or
internals and do not need the read wrapped in Effect.
Details
After endUnsafe, a queue remains Closing while buffered messages are
drained, and its size continues to include those messages. A Done queue
reports a size of 0. This unsafe operation reads the queue state directly
without Effect wrapping.
Signature
declare function sizeUnsafe<A, E>(self: Dequeue<A, E>): numberExample
(Checking queue size synchronously)
import { Cause, Effect, Queue } from "effect"
const program = Effect.gen(function*() { const queue = yield* Queue.bounded<number, Cause.Done>(10)
// Check size of empty queue const size1 = Queue.sizeUnsafe(queue)
// Add some messages Queue.offerUnsafe(queue, 1) Queue.offerUnsafe(queue, 2) Queue.offerUnsafe(queue, 3)
// Check size after adding messages const size2 = Queue.sizeUnsafe(queue)
// End the queue Queue.endUnsafe(queue)
// Ending retains the buffered size while the queue is Closing const size3 = Queue.sizeUnsafe(queue) return [size1, size2, size3]})
await Effect.runPromise(program) // => [0, 3, 3]Taking
Takes and returns all currently buffered messages without waiting for more.
Details
Returns an empty array when the queue is empty or has completed normally. If the queue has failed, the effect fails with the queue's error.
Signature
declare function clear<A, E>(self: Dequeue<A, E>): Effect<Array<A>, Exclude<E, Done<any>>>Example
(Clearing queued values)
import { Effect, Queue } from "effect"
const program = Effect.gen(function*() { const queue = yield* Queue.bounded<number>(10)
// Add several messages yield* Queue.offerAll(queue, [1, 2, 3, 4, 5])
// Clear all messages from the queue const messages = yield* Queue.clear(queue)
// Queue is now empty const size = yield* Queue.size(queue)
// Clearing empty queue returns empty array const empty = yield* Queue.clear(queue) return { messages, size, empty }})
await Effect.runPromise(program) // => { messages: [1, 2, 3, 4, 5], size: 0, empty: [] }Takes all messages from the queue, until the queue has errored or is done.
Signature
declare function collect<A, E>(self: Dequeue<A, Done<void> | E>): Effect<Array<A>, Exclude<E, Done<any>>>Example
(Collecting values until completion)
import { Cause, Effect, Queue } from "effect"
const program = Effect.gen(function*() { const queue = yield* Queue.bounded<number, Cause.Done>(5)
// Add several messages yield* Queue.offerAll(queue, [1, 2, 3, 4, 5]) yield* Queue.end(queue)
// Collect all available messages return yield* Queue.collect(queue)})
await Effect.runPromise(program) // => [1, 2, 3, 4, 5]Peeks at the next item without removing it.
Details
Blocks until an item is available. If the queue is done or fails, the error is propagated.
Signature
declare function peek<A, E>(self: Dequeue<A, E>): Effect<A, E>Example
(Peeking at the next value)
import { Cause, Effect, Queue } from "effect"
const program = Effect.gen(function*() { const queue = yield* Queue.bounded<number>(10) yield* Queue.offer(queue, 42)
// Peek at the next item without removing it const item = yield* Queue.peek(queue) return item})
await Effect.runPromise(program) // => 42Attempts to take one item from the queue without waiting.
Details
Returns Option.some when an item is immediately available. Returns
Option.none when no item is available, when the queue is done, or when the
immediate take observes a queue failure.
Signature
declare function poll<A, E>(self: Dequeue<A, E>): Effect<Option<A>>Example
(Polling without blocking)
import { Effect, Option, Queue } from "effect"
const program = Effect.gen(function*() { const queue = yield* Queue.bounded<number>(10)
// Poll returns Option.none if empty const maybe1 = yield* Queue.poll(queue)
// Add an item yield* Queue.offer(queue, 42)
// Poll returns Option.some with the item const maybe2 = yield* Queue.poll(queue) return [maybe1, maybe2]})
await Effect.runPromise(program) // => [Option.none(), Option.some(42)]Takes a single message from the queue, or wait for a message to be available.
Details
If the queue is done, it will fail with Done. If the
queue fails, the Effect will fail with the error.
Signature
declare function take<A, E>(self: Dequeue<A, E>): Effect<A, E>Example
(Taking one value)
import { Cause, Effect, Exit, Queue } from "effect"
const program = Effect.gen(function*() { const queue = yield* Queue.bounded<string, Cause.Done>(3)
// Add some messages yield* Queue.offer(queue, "first") yield* Queue.offer(queue, "second")
// Take messages one by one const msg1 = yield* Queue.take(queue) const msg2 = yield* Queue.take(queue)
// End the queue yield* Queue.end(queue)
// Taking from an ended queue fails with Done const result = yield* Effect.exit(Queue.take(queue)) return [[msg1, msg2], result]})
await Effect.runPromise(program) // => [["first", "second"], Exit.fail(Cause.Done())]Takes all currently available messages, waiting until at least one message is available when the queue is empty.
When to use
Use when consumers should process the next non-empty batch of buffered messages instead of repeatedly taking one message at a time.
Details
Returns a non-empty array. If the queue completes or fails before a message can be taken, the effect fails with the queue's terminal error.
Signature
declare function takeAll<A, E>(self: Dequeue<A, E>): Effect<[A, ...Array<A>], E>Example
(Taking all available values)
import { Cause, Effect, Queue } from "effect"
const program = Effect.gen(function*() { const queue = yield* Queue.bounded<number, Cause.Done>(5)
// Add several messages yield* Queue.offerAll(queue, [1, 2, 3, 4, 5])
// Take all available messages const messages1 = yield* Queue.takeAll(queue) return messages1})
await Effect.runPromise(program) // => [1, 2, 3, 4, 5]takeBetween
Takes between min and max messages from the queue.
Details
The operation waits when fewer than the required minimum messages are
available. It returns at most max messages. If the queue completes or fails
before the minimum can be satisfied, the effect fails with the queue's
terminal error.
Signature
declare function takeBetween<A, E>(self: Dequeue<A, E>, min: number, max: number): Effect<Array<A>, E>Example
(Taking a bounded batch of values)
import { Cause, Effect, Queue } from "effect"
const program = Effect.gen(function*() { const queue = yield* Queue.bounded<number>(10)
// Add several messages yield* Queue.offerAll(queue, [1, 2, 3, 4, 5, 6, 7, 8])
// Take between 2 and 5 messages const batch1 = yield* Queue.takeBetween(queue, 2, 5)
// Take between 1 and 10 messages (but only 3 remain) const batch2 = yield* Queue.takeBetween(queue, 1, 10)
// No more messages available, will wait or return done // const batch3 = yield* Queue.takeBetween(queue, 1, 3) return [batch1, batch2]})
await Effect.runPromise(program) // => [[1, 2, 3, 4, 5], [6, 7, 8]]Takes up to n messages from the queue.
Details
The operation may wait until enough messages are available to satisfy the
queue's batching rules. If n is less than or equal to zero, it succeeds
with an empty array. If the queue completes or fails before messages can be
taken, the effect fails with the queue's terminal error.
Signature
declare function takeN<A, E>(self: Dequeue<A, E>, n: number): Effect<Array<A>, E>Example
(Taking a fixed number of values)
import { Cause, Effect, Queue } from "effect"
const program = Effect.gen(function*() { const queue = yield* Queue.bounded<number, Cause.Done>(10)
// Add several messages yield* Queue.offerAll(queue, [1, 2, 3, 4, 5, 6, 7])
// Take exactly 3 messages const first3 = yield* Queue.takeN(queue, 3)
// Take exactly 2 more messages const next2 = yield* Queue.takeN(queue, 2)
// Take remaining messages const remaining = yield* Queue.takeN(queue, 2) return [first3, next2, remaining]})
await Effect.runPromise(program) // => [[1, 2, 3], [4, 5], [6, 7]]takeUnsafe
Attempts to take one message from the queue synchronously.
When to use
Use when polling queue internals must not suspend or register a waiting taker,
and undefined is an acceptable result for an empty queue.
Details
Returns an Exit for an immediately available message or for the queue's
terminal state. Returns undefined when no message is immediately available.
This operation does not wait or register a taker.
Signature
declare function takeUnsafe<A, E>(self: Dequeue<A, E>): Exit<A, E> | undefinedExample
(Taking one value synchronously)
import { Effect, Exit, Queue } from "effect"
// Create a queue and use unsafe operationsconst program = Effect.gen(function*() { const queue = yield* Queue.bounded<number>(10)
// Add some messages Queue.offerUnsafe(queue, 1) Queue.offerUnsafe(queue, 2)
// Take a message synchronously const result1 = Queue.takeUnsafe(queue)
const result2 = Queue.takeUnsafe(queue)
// No more messages - returns undefined const result3 = Queue.takeUnsafe(queue) return [result1, result2, result3]})
await Effect.runPromise(program) // => [Exit.succeed(1), Exit.succeed(2), undefined]