(self: TxQueueState): Effect.Effect<void>Waits for the queue to complete (either successfully or with failure).
Example (Awaiting queue completion)
import { Effect, TxQueue } from "effect"
const program = Effect.gen(function*() {
const queue = yield* TxQueue.bounded<number, string>(10)
// In another fiber, end the queue
yield* Effect.forkChild(Effect.delay(TxQueue.interrupt(queue), "100 millis"))
// Wait for completion - succeeds when queue ends
yield* TxQueue.awaitCompletion(queue)
console.log("Queue completed successfully")
})combinators
Source effect/TxQueue.ts:147811 lines
export const const awaitCompletion: (
self: TxQueueState
) => Effect.Effect<void>
Waits for the queue to complete (either successfully or with failure).
Example (Awaiting queue completion)
import { Effect, TxQueue } from "effect"
const program = Effect.gen(function*() {
const queue = yield* TxQueue.bounded<number, string>(10)
// In another fiber, end the queue
yield* Effect.forkChild(Effect.delay(TxQueue.interrupt(queue), "100 millis"))
// Wait for completion - succeeds when queue ends
yield* TxQueue.awaitCompletion(queue)
console.log("Queue completed successfully")
})
awaitCompletion = (self: TxQueueState(parameter) self: {
strategy: "bounded" | "unbounded" | "dropping" | "sliding";
capacity: number;
items: TxChunk.TxChunk<any>;
stateRef: TxRef.TxRef<State<any, any>>;
toString: () => string;
toJSON: () => unknown;
}
self: TxQueueState): import EffectEffect.type Effect.Effect = /*unresolved*/ anyEffect<void> =>
import EffectEffect.gen(function*() {
const const state: State<any, any>state = yield* import TxRefTxRef.const get: <A>(
self: TxRef<A>
) => Effect.Effect<A>
Reads the current value of the TxRef.
When to use
Use to read the current value of a TxRef.
Example (Reading transactional references)
import { Effect, TxRef } from "effect"
const program = Effect.gen(function*() {
const counter = yield* TxRef.make(42)
// Read the value within a transaction
const value = yield* Effect.tx(
TxRef.get(counter)
)
console.log(value) // 42
})
get(self: TxQueueState(parameter) self: {
strategy: "bounded" | "unbounded" | "dropping" | "sliding";
capacity: number;
items: TxChunk.TxChunk<any>;
stateRef: TxRef.TxRef<State<any, any>>;
toString: () => string;
toJSON: () => unknown;
}
self.TxQueueState.stateRef: TxRef.TxRef<State<any, any>>(property) TxQueueState.stateRef: {
version: number;
pending: Map<unknown, () => void>;
value: A;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
stateRef)
if (const state: State<any, any>state._tag === "Done") {
return void 0
}
// Not done yet, retry transaction
return yield* import EffectEffect.txRetry
}).pipe(import EffectEffect.tx)