Skip to content
Effect Days 2026 Get your ticket

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.

42 exports Added in v2.0.0 Source

Completion

end

Added in v4.0.0 Source

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]

endUnsafe

Added in v4.0.0 Source

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>): boolean

Example

(Ending queues synchronously)

import { Cause, Effect, Queue } from "effect"
// Create a queue and use unsafe operations
const 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"] }

fail

Added in v4.0.0 Source

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")]

failCause

Added in v4.0.0 Source

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"))]

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>): boolean

Example

(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"))]

interrupt

Added in v4.0.0 Source

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 }

into

Added in v4.0.0 Source

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())]

shutdown

Added in v2.0.0 Source

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

bounded

Added in v2.0.0 Source

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) // => 2

dropping

Added in v2.0.0 Source

Creates 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]]

make

Added in v4.0.0 Source

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 }

sliding

Added in v2.0.0 Source

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]

unbounded

Added in v2.0.0 Source

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

asDequeue

Added in v4.0.0 Source

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

  • asEnqueue for narrowing a queue to its producer side
  • Dequeue for the consumer-side queue handle returned by this function

Signature

declare const asDequeue: <A, E>(self: Queue<A, E>) => Dequeue<A, E>

asEnqueue

Added in v4.0.0 Source

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

  • asDequeue for exposing only the read side of a Queue
  • Enqueue for the write-only queue handle returned by this conversion

Signature

declare function asEnqueue<A, E>(self: Queue<A, E>): Enqueue<A, E>

Guards

isDequeue

Added in v2.0.0 Source

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

  • Dequeue for the read-side queue handle checked by this guard
  • isQueue for checking for a full read-write queue handle
  • isEnqueue for checking for the write side of a queue
  • asDequeue for narrowing an existing Queue to its read-only interface

Signature

declare function isDequeue<A = unknown, E = unknown>(u: unknown): u is Dequeue<A, E>

isEnqueue

Added in v2.0.0 Source

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

  • isQueue for checking for a full read-write queue handle
  • isDequeue for checking for the read side of a queue
  • asEnqueue for narrowing an existing Queue to its write-only interface

Signature

declare function isEnqueue<A = unknown, E = unknown>(u: unknown): u is Enqueue<A, E>

isQueue

Added in v2.0.0 Source

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

  • isEnqueue for checking values that only need write access
  • isDequeue for checking values that only need read access

Signature

declare function isQueue<A = unknown, E = unknown>(u: unknown): u is Queue<A, E>

Models

Dequeue interface

Added in v2.0.0 Source

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"

Enqueue interface

Added in v2.0.0 Source

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 queue
const 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", "!"]

Queue interface

Added in v2.0.0 Source

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

offer

Added in v2.0.0 Source

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 }

offerAll

Added in v2.0.0 Source

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]

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 API
const 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

Added in v4.0.0 Source

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>): boolean

Example

(Offering a value synchronously)

import { Effect, Queue } from "effect"
// Create a queue effect and extract the queue for unsafe operations
const 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>>>

Dequeue

Added in v2.0.0 Source

Companion namespace containing type-level metadata for the Dequeue read-only queue interface.

Enqueue

Added in v2.0.0 Source

Companion namespace containing type-level metadata for the Enqueue write-only queue interface.

Queue

Added in v2.0.0 Source

Companion namespace containing type-level metadata and low-level state types for Queue.

Predicates

isFull

Added in v2.0.0 Source

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

Added in v4.0.0 Source

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>): boolean

Example

(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

size

Added in v2.0.0 Source

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

Added in v4.0.0 Source

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>): number

Example

(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

clear

Added in v4.0.0 Source

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: [] }

collect

Added in v4.0.0 Source

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]

peek

Added in v4.0.0 Source

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) // => 42

poll

Added in v2.0.0 Source

Attempts 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)]

take

Added in v2.0.0 Source

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())]

takeAll

Added in v2.0.0 Source

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

Added in v2.0.0 Source

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]]

takeN

Added in v2.0.0 Source

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

Added in v4.0.0 Source

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> | undefined

Example

(Taking one value synchronously)

import { Effect, Exit, Queue } from "effect"
// Create a queue and use unsafe operations
const 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]