Skip to content
Effect Days 2026 Get your ticket

TxPriorityQueue

Transactional priority queues whose state is stored in a TxRef. Elements are kept in the order defined by the Order supplied at construction time, and dequeue operations return the first element according to that ordering.

Use TxPriorityQueue when multiple fibers coordinate through a shared queue and queue operations need to compose with other transactional state changes. The retrying peek and take operations wait transactionally when the queue is empty, so they can be combined with other transactional reads and writes in one atomic workflow.

19 exports Added in v2.0.0 Source

Constructors

empty

Added in v2.0.0 Source

Creates an empty TxPriorityQueue with the given ordering.

Signature

declare function empty<A>(order: Order<A>): Effect<TxPriorityQueue<A>>

Example

(Creating an empty priority queue)

import { Effect, Order, TxPriorityQueue } from "effect"
const program = Effect.gen(function*() {
const pq = yield* TxPriorityQueue.empty<number>(Order.Number)
return yield* TxPriorityQueue.isEmpty(pq)
})
await Effect.runPromise(program) // => true

fromIterable

Added in v2.0.0 Source

Creates a TxPriorityQueue from an iterable of elements.

Signature

declare const fromIterable: {
<A>(order: Order<A>): (iterable: Iterable<A>) => Effect<TxPriorityQueue<A>>;
<A>(order: Order<A>, iterable: Iterable<A>): Effect<TxPriorityQueue<A>>;
}

Example

(Creating a priority queue from an iterable)

import { Effect, Order, TxPriorityQueue } from "effect"
const program = Effect.gen(function*() {
const pq = yield* TxPriorityQueue.fromIterable(Order.Number, [3, 1, 2])
return yield* TxPriorityQueue.take(pq)
})
await Effect.runPromise(program) // => 1

make

Added in v2.0.0 Source

Creates a TxPriorityQueue from variadic elements.

Signature

declare function make<A>(order: Order<A>): (...elements: Array<A>) => Effect<TxPriorityQueue<A>>

Example

(Creating a priority queue from variadic values)

import { Effect, Order, TxPriorityQueue } from "effect"
const program = Effect.gen(function*() {
const pq = yield* TxPriorityQueue.make(Order.Number)(3, 1, 2)
return yield* TxPriorityQueue.take(pq)
})
await Effect.runPromise(program) // => 1

Converting

toArray

Added in v2.0.0 Source

Returns all elements in priority order without removing them.

Signature

declare function toArray<A>(self: TxPriorityQueue<A>): Effect<Array<A>>

Example

(Reading values in priority order)

import { Effect, Order, TxPriorityQueue } from "effect"
const program = Effect.gen(function*() {
const pq = yield* TxPriorityQueue.fromIterable(Order.Number, [3, 1, 2])
return yield* TxPriorityQueue.toArray(pq)
})
await Effect.runPromise(program) // => [1, 2, 3]

Filtering

removeIf

Added in v2.0.0 Source

Removes elements matching the predicate.

Signature

declare const removeIf: {
<A>(predicate: Predicate<A>): (self: TxPriorityQueue<A>) => Effect<void>;
<A>(self: TxPriorityQueue<A>, predicate: Predicate<A>): Effect<void>;
}

Example

(Removing matching values)

import { Effect, Order, TxPriorityQueue } from "effect"
const program = Effect.gen(function*() {
const pq = yield* TxPriorityQueue.fromIterable(Order.Number, [1, 2, 3, 4, 5])
yield* TxPriorityQueue.removeIf(pq, (n) => n % 2 === 0)
return yield* TxPriorityQueue.takeAll(pq)
})
await Effect.runPromise(program) // => [1, 3, 5]

retainIf

Added in v2.0.0 Source

Keeps only elements matching the predicate.

Signature

declare const retainIf: {
<A>(predicate: Predicate<A>): (self: TxPriorityQueue<A>) => Effect<void>;
<A>(self: TxPriorityQueue<A>, predicate: Predicate<A>): Effect<void>;
}

Example

(Retaining matching values)

import { Effect, Order, TxPriorityQueue } from "effect"
const program = Effect.gen(function*() {
const pq = yield* TxPriorityQueue.fromIterable(Order.Number, [1, 2, 3, 4, 5])
yield* TxPriorityQueue.retainIf(pq, (n) => n % 2 === 0)
return yield* TxPriorityQueue.takeAll(pq)
})
await Effect.runPromise(program) // => [2, 4]

Getters

peek

Added in v2.0.0 Source

Observes the smallest element without removing it.

When to use

Use to inspect the next prioritized value and retry transactionally while the queue is empty.

Signature

declare function peek<A>(self: TxPriorityQueue<A>): Effect<A>

Example

(Peeking at the next value)

import { Effect, Order, TxPriorityQueue } from "effect"
const program = Effect.gen(function*() {
const pq = yield* TxPriorityQueue.fromIterable(Order.Number, [3, 1, 2])
return yield* TxPriorityQueue.peek(pq)
})
await Effect.runPromise(program) // => 1

peekOption

Added in v2.0.0 Source

Observes the smallest element without removing it, returning None when the queue is empty.

When to use

Use to inspect the next prioritized value without retrying on an empty queue.

Signature

declare function peekOption<A>(self: TxPriorityQueue<A>): Effect<Option<A>>

Example

(Peeking without retrying)

import { Effect, Option, Order, TxPriorityQueue } from "effect"
const program = Effect.gen(function*() {
const pq = yield* TxPriorityQueue.empty<number>(Order.Number)
return yield* TxPriorityQueue.peekOption(pq)
})
await Effect.runPromise(program) // => Option.none()

size

Added in v2.0.0 Source

Returns the number of elements in the queue.

Signature

declare function size<A>(self: TxPriorityQueue<A>): Effect<number>

Example

(Getting the queue size)

import { Effect, Order, TxPriorityQueue } from "effect"
const program = Effect.gen(function*() {
const pq = yield* TxPriorityQueue.fromIterable(Order.Number, [1, 2, 3])
return yield* TxPriorityQueue.size(pq)
})
await Effect.runPromise(program) // => 3

Guards

Determines if the provided value is a TxPriorityQueue.

Signature

declare function isTxPriorityQueue(u: unknown): u is TxPriorityQueue<unknown>

Example

(Checking for a TxPriorityQueue)

import { Effect, Order, TxPriorityQueue } from "effect"
const program = Effect.gen(function*() {
const pq = yield* TxPriorityQueue.empty<number>(Order.Number)
return [TxPriorityQueue.isTxPriorityQueue(pq), TxPriorityQueue.isTxPriorityQueue("nope")]
})
await Effect.runPromise(program) // => [true, false]

Models

TxPriorityQueue interface

Added in v4.0.0 Source

A transactional priority queue backed by a sorted Chunk.

Details

Elements are stored in ascending order according to the Order provided at construction time. take returns the smallest element, peek observes it without removing.

Signature

interface TxPriorityQueue<in out A> extends Inspectable, Pipeable {
readonly "~effect/transactions/TxPriorityQueue": "~effect/transactions/TxPriorityQueue";
readonly ord: Order<A>;
readonly ref: TxRef<Chunk<A>>;
}

Example

(Dequeuing values by priority)

import { Effect, Order, TxPriorityQueue } from "effect"
const program = Effect.gen(function*() {
const pq = yield* TxPriorityQueue.empty<number>(Order.Number)
yield* TxPriorityQueue.offer(pq, 3)
yield* TxPriorityQueue.offer(pq, 1)
yield* TxPriorityQueue.offer(pq, 2)
return yield* TxPriorityQueue.take(pq)
})
await Effect.runPromise(program) // => 1

Mutations

offer

Added in v2.0.0 Source

Inserts an element into the queue in sorted position.

Signature

declare const offer: {
<A>(value: A): (self: TxPriorityQueue<A>) => Effect<void>;
<A>(self: TxPriorityQueue<A>, value: A): Effect<void>;
}

Example

(Offering a value)

import { Effect, Order, TxPriorityQueue } from "effect"
const program = Effect.gen(function*() {
const pq = yield* TxPriorityQueue.empty<number>(Order.Number)
yield* TxPriorityQueue.offer(pq, 2)
yield* TxPriorityQueue.offer(pq, 1)
return yield* TxPriorityQueue.take(pq)
})
await Effect.runPromise(program) // => 1

offerAll

Added in v2.0.0 Source

Inserts all elements from an iterable into the queue.

Signature

declare const offerAll: {
<A>(values: Iterable<A>): (self: TxPriorityQueue<A>) => Effect<void>;
<A>(self: TxPriorityQueue<A>, values: Iterable<A>): Effect<void>;
}

Example

(Offering multiple values)

import { Effect, Order, TxPriorityQueue } from "effect"
const program = Effect.gen(function*() {
const pq = yield* TxPriorityQueue.empty<number>(Order.Number)
yield* TxPriorityQueue.offerAll(pq, [3, 1, 2])
return yield* TxPriorityQueue.take(pq)
})
await Effect.runPromise(program) // => 1

take

Added in v2.0.0 Source

Takes the smallest element from the queue. Retries if the queue is empty.

Signature

declare function take<A>(self: TxPriorityQueue<A>): Effect<A>

Example

(Taking the next value)

import { Effect, Order, TxPriorityQueue } from "effect"
const program = Effect.gen(function*() {
const pq = yield* TxPriorityQueue.fromIterable(Order.Number, [3, 1, 2])
return yield* TxPriorityQueue.take(pq)
})
await Effect.runPromise(program) // => 1

takeAll

Added in v2.0.0 Source

Takes all elements from the queue, returning them in priority order.

Signature

declare function takeAll<A>(self: TxPriorityQueue<A>): Effect<Array<A>>

Example

(Taking all values in priority order)

import { Effect, Order, TxPriorityQueue } from "effect"
const program = Effect.gen(function*() {
const pq = yield* TxPriorityQueue.fromIterable(Order.Number, [3, 1, 2])
return yield* TxPriorityQueue.takeAll(pq)
})
await Effect.runPromise(program) // => [1, 2, 3]

takeOption

Added in v2.0.0 Source

Tries to take the smallest element. Returns None if the queue is empty.

Signature

declare function takeOption<A>(self: TxPriorityQueue<A>): Effect<Option<A>>

Example

(Taking without retrying)

import { Effect, Option, Order, TxPriorityQueue } from "effect"
const program = Effect.gen(function*() {
const pq = yield* TxPriorityQueue.empty<number>(Order.Number)
return yield* TxPriorityQueue.takeOption(pq)
})
await Effect.runPromise(program) // => Option.none()

takeUpTo

Added in v2.0.0 Source

Takes up to n elements from the queue in priority order.

Signature

declare const takeUpTo: {
(n: number): <A>(self: TxPriorityQueue<A>) => Effect<Array<A>>;
<A>(self: TxPriorityQueue<A>, n: number): Effect<Array<A>>;
}

Example

(Taking up to a limit)

import { Effect, Order, TxPriorityQueue } from "effect"
const program = Effect.gen(function*() {
const pq = yield* TxPriorityQueue.fromIterable(Order.Number, [5, 3, 1, 4, 2])
return yield* TxPriorityQueue.takeUpTo(pq, 2)
})
await Effect.runPromise(program) // => [1, 2]

Predicates

isEmpty

Added in v2.0.0 Source

Returns true if the queue is empty.

Signature

declare function isEmpty<A>(self: TxPriorityQueue<A>): Effect<boolean>

Example

(Checking whether a queue is empty)

import { Effect, Order, TxPriorityQueue } from "effect"
const program = Effect.gen(function*() {
const pq = yield* TxPriorityQueue.empty<number>(Order.Number)
return yield* TxPriorityQueue.isEmpty(pq)
})
await Effect.runPromise(program) // => true

isNonEmpty

Added in v2.0.0 Source

Returns true if the queue has at least one element.

Signature

declare function isNonEmpty<A>(self: TxPriorityQueue<A>): Effect<boolean>

Example

(Checking whether a queue has elements)

import { Effect, Order, TxPriorityQueue } from "effect"
const program = Effect.gen(function*() {
const pq = yield* TxPriorityQueue.fromIterable(Order.Number, [1])
return yield* TxPriorityQueue.isNonEmpty(pq)
})
await Effect.runPromise(program) // => true