<A, E>(self: Enqueue<A, E | Done>): booleanSignals 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.
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)
console.log(ended) // true
// Existing messages can still be consumed while the queue is closing
console.log(queue.state._tag) // "Closing"
Queue.takeUnsafe(queue)
Queue.takeUnsafe(queue)
// After buffered messages are consumed, the queue is done
console.log(queue.state._tag) // "Done"
})export const const endUnsafe: <A, E>(
self: Enqueue<A, E | Done>
) => boolean
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.
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)
console.log(ended) // true
// Existing messages can still be consumed while the queue is closing
console.log(queue.state._tag) // "Closing"
Queue.takeUnsafe(queue)
Queue.takeUnsafe(queue)
// After buffered messages are consumed, the queue is done
console.log(queue.state._tag) // "Done"
})
endUnsafe = <function (type parameter) A in <A, E>(self: Enqueue<A, E | Done>): booleanA, function (type parameter) E in <A, E>(self: Enqueue<A, E | Done>): booleanE>(self: Enqueue<A, E | Done>(parameter) self: {
strategy: "suspend" | "dropping" | "sliding";
dispatcher: SchedulerDispatcher;
capacity: number;
messages: MutableList.MutableList<any>;
state: Queue.State<any, any>;
scheduleRunning: boolean;
toString: () => string;
toJSON: () => unknown;
}
self: interface Enqueue<in A, in E = never>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.
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)
})
Companion namespace containing type-level metadata for the Enqueue
write-only queue interface.
Enqueue<function (type parameter) A in <A, E>(self: Enqueue<A, E | Done>): booleanA, function (type parameter) E in <A, E>(self: Enqueue<A, E | Done>): booleanE | import DoneDone>) => const failCauseUnsafe: <A, E>(
self: Enqueue<A, E>,
cause: Cause<E>
) => boolean
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.
Example (Failing queues with a cause synchronously)
import { Cause, Effect, 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)
console.log(failed) // true
// The queue is now done with the specified failure cause
console.log(queue.state._tag) // "Done"
})
failCauseUnsafe(self: Enqueue<A, E | Done>(parameter) self: {
strategy: "suspend" | "dropping" | "sliding";
dispatcher: SchedulerDispatcher;
capacity: number;
messages: MutableList.MutableList<any>;
state: Queue.State<any, any>;
scheduleRunning: boolean;
toString: () => string;
toJSON: () => unknown;
}
self, import corecore.const causeFail: <E>(
error: E
) => Cause.Cause<E>
causeFail(import corecore.const Done: <A = void>(
value?: A
) => Cause.Done<A>
Done()))