<A>(subscription: PubSub.Subscription<A>): Channel<
Arr.NonEmptyReadonlyArray<A>
>Creates a channel from a PubSub subscription that outputs arrays of values.
Details
This constructor creates a channel that reads from a PubSub subscription and outputs arrays of values in chunks. It's useful when you want to process multiple values at once for better performance.
Example (Batching subscription values)
import { Channel, Data, Effect, PubSub } from "effect"
class StreamError extends Data.TaggedError("StreamError")<{
readonly message: string
}> {}
const program = Effect.gen(function*() {
const pubsub = yield* PubSub.bounded<number>(16)
const subscription = yield* PubSub.subscribe(pubsub)
// Create a channel that reads arrays of values
const channel = Channel.fromSubscriptionArray(subscription)
// Publish some values
yield* PubSub.publish(pubsub, 1)
yield* PubSub.publish(pubsub, 2)
yield* PubSub.publish(pubsub, 3)
yield* PubSub.publish(pubsub, 4)
// The channel will output arrays like [1, 2, 3] and [4]
return channel
})Example (Processing subscription values in batches)
import { Channel, Data, Effect, PubSub } from "effect"
class BatchProcessingError extends Data.TaggedError("BatchProcessingError")<{
readonly reason: string
}> {}
const batchProcessor = Effect.gen(function*() {
const pubsub = yield* PubSub.bounded<string>(32)
const subscription = yield* PubSub.subscribe(pubsub)
// Create a channel that processes items in batches
const batchChannel = Channel.fromSubscriptionArray(subscription)
// Transform to process each batch
const processedChannel = Channel.map(batchChannel, (batch) => {
console.log(`Processing batch of ${batch.length} items:`, batch)
return batch.map((item) => item.toUpperCase())
})
return processedChannel
})Example (Aggregating subscription metrics)
import { Channel, Effect, PubSub } from "effect"
const metricsAggregator = Effect.gen(function*() {
const metricsPubSub = yield* PubSub.bounded<
{ timestamp: number; value: number }
>(100)
const subscription = yield* PubSub.subscribe(metricsPubSub)
// Create a channel that collects metrics in chunks
const metricsChannel = Channel.fromSubscriptionArray(subscription)
// Transform to calculate aggregate statistics
const aggregatedChannel = Channel.map(metricsChannel, (metrics) => {
const values = metrics.map((m) => m.value)
const sum = values.reduce((a, b) => a + b, 0)
const avg = sum / values.length
const min = Math.min(...values)
const max = Math.max(...values)
return {
count: values.length,
sum,
average: avg,
min,
max,
firstTimestamp: Math.min(...metrics.map((m) => m.timestamp)),
lastTimestamp: Math.max(...metrics.map((m) => m.timestamp))
}
})
return aggregatedChannel
})export const const fromSubscriptionArray: <A>(
subscription: PubSub.Subscription<A>
) => Channel<Arr.NonEmptyReadonlyArray<A>>
Creates a channel from a PubSub subscription that outputs arrays of values.
Details
This constructor creates a channel that reads from a PubSub subscription and outputs
arrays of values in chunks. It's useful when you want to process multiple values at once
for better performance.
Example (Batching subscription values)
import { Channel, Data, Effect, PubSub } from "effect"
class StreamError extends Data.TaggedError("StreamError")<{
readonly message: string
}> {}
const program = Effect.gen(function*() {
const pubsub = yield* PubSub.bounded<number>(16)
const subscription = yield* PubSub.subscribe(pubsub)
// Create a channel that reads arrays of values
const channel = Channel.fromSubscriptionArray(subscription)
// Publish some values
yield* PubSub.publish(pubsub, 1)
yield* PubSub.publish(pubsub, 2)
yield* PubSub.publish(pubsub, 3)
yield* PubSub.publish(pubsub, 4)
// The channel will output arrays like [1, 2, 3] and [4]
return channel
})
Example (Processing subscription values in batches)
import { Channel, Data, Effect, PubSub } from "effect"
class BatchProcessingError extends Data.TaggedError("BatchProcessingError")<{
readonly reason: string
}> {}
const batchProcessor = Effect.gen(function*() {
const pubsub = yield* PubSub.bounded<string>(32)
const subscription = yield* PubSub.subscribe(pubsub)
// Create a channel that processes items in batches
const batchChannel = Channel.fromSubscriptionArray(subscription)
// Transform to process each batch
const processedChannel = Channel.map(batchChannel, (batch) => {
console.log(`Processing batch of ${batch.length} items:`, batch)
return batch.map((item) => item.toUpperCase())
})
return processedChannel
})
Example (Aggregating subscription metrics)
import { Channel, Effect, PubSub } from "effect"
const metricsAggregator = Effect.gen(function*() {
const metricsPubSub = yield* PubSub.bounded<
{ timestamp: number; value: number }
>(100)
const subscription = yield* PubSub.subscribe(metricsPubSub)
// Create a channel that collects metrics in chunks
const metricsChannel = Channel.fromSubscriptionArray(subscription)
// Transform to calculate aggregate statistics
const aggregatedChannel = Channel.map(metricsChannel, (metrics) => {
const values = metrics.map((m) => m.value)
const sum = values.reduce((a, b) => a + b, 0)
const avg = sum / values.length
const min = Math.min(...values)
const max = Math.max(...values)
return {
count: values.length,
sum,
average: avg,
min,
max,
firstTimestamp: Math.min(...metrics.map((m) => m.timestamp)),
lastTimestamp: Math.max(...metrics.map((m) => m.timestamp))
}
})
return aggregatedChannel
})
fromSubscriptionArray = <function (type parameter) A in <A>(subscription: PubSub.Subscription<A>): Channel<Arr.NonEmptyReadonlyArray<A>>A>(
subscription: PubSub.Subscription<A>(parameter) subscription: {
pubsub: PubSub.Atomic<any>;
subscribers: PubSub.Subscribers<any>;
subscription: PubSub.BackingSubscription<A>;
pollers: MutableList.MutableList<Deferred.Deferred<any>>;
shutdownHook: Latch.Latch;
shutdownFlag: MutableRef.MutableRef<boolean>;
strategy: PubSub.Strategy<any>;
replayWindow: PubSub.ReplayWindow<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; <…;
}
subscription: import PubSubPubSub.interface Subscription<out A>A subscription represents a consumer's connection to a PubSub, allowing them to take messages.
Example (Taking messages from a subscription)
import { Effect, PubSub } from "effect"
const program = Effect.gen(function*() {
const pubsub = yield* PubSub.bounded<string>(10)
// Subscribe within a scope for automatic cleanup
yield* Effect.scoped(Effect.gen(function*() {
const subscription: PubSub.Subscription<string> = yield* PubSub.subscribe(
pubsub
)
yield* PubSub.publishAll(pubsub, ["msg1", "msg2", "msg3"])
// Take individual messages
const message = yield* PubSub.take(subscription)
console.log(message) // "msg1"
// Take multiple messages
const messages = yield* PubSub.takeUpTo(subscription, 1)
console.log(messages) // ["msg2"]
const allMessages = yield* PubSub.takeAll(subscription)
console.log(allMessages) // ["msg3"]
}))
})
Subscription<function (type parameter) A in <A>(subscription: PubSub.Subscription<A>): Channel<Arr.NonEmptyReadonlyArray<A>>A>
): interface Channel<out OutElem, out OutErr = never, out OutDone = void, in InElem = unknown, in InErr = unknown, in InDone = unknown, out Env = never>A Channel is a nexus of I/O operations, which supports both reading and
writing. A channel may read values of type InElem and write values of type
OutElem. When the channel finishes, it yields a value of type OutDone. A
channel may fail with a value of type OutErr.
Details
Channels are the foundation of Streams: both streams and sinks are built on
channels. Most users shouldn't have to use channels directly, as streams and
sinks are much more convenient and cover all common use cases. However, when
adding new stream and sink operators, or doing something highly specialized,
it may be useful to use channels directly.
Channels compose in a variety of ways:
- Piping: One channel can be piped to another channel, assuming the
input type of the second is the same as the output type of the first.
- Sequencing: The terminal value of one channel can be used to create
another channel, and both the first channel and the function that makes
the second channel can be composed into a channel.
- Concatenating: The output of one channel can be used to create other
channels, which are all concatenated together. The first channel and the
function that makes the other channels can be composed into a channel.
Example (Typing channels)
import type { Channel } from "effect"
// A channel that outputs numbers and requires no environment
type NumberChannel = Channel.Channel<number>
// A channel that outputs strings, can fail with Error, completes with boolean
type StringChannel = Channel.Channel<string, Error, boolean>
// A channel with all type parameters specified
type FullChannel = Channel.Channel<
string, // OutElem - output elements
Error, // OutErr - output errors
number, // OutDone - completion value
number, // InElem - input elements
string, // InErr - input errors
boolean, // InDone - input completion
{ db: string } // Env - required environment
>
Channel<import ArrArr.type Arr.NonEmptyReadonlyArray = /*unresolved*/ anyNonEmptyReadonlyArray<function (type parameter) A in <A>(subscription: PubSub.Subscription<A>): Channel<Arr.NonEmptyReadonlyArray<A>>A>> =>
const fromPull: <
OutElem,
OutErr,
OutDone,
EX,
EnvX,
Env
>(
effect: Effect.Effect<
Pull.Pull<OutElem, OutErr, OutDone, EnvX>,
EX,
Env
>
) => Channel<
OutElem,
Pull.ExcludeDone<OutErr> | EX,
OutDone,
unknown,
unknown,
unknown,
Env | EnvX
>
Creates a Channel from an Effect that produces a Pull.
Example (Creating channels from pulls)
import { Channel, Effect } from "effect"
const channel = Channel.fromPull(
Effect.succeed(Effect.succeed(42))
)
fromPull(import EffectEffect.const succeed: <A>(value: A) => Effect<A>Creates an Effect that always succeeds with a given value.
When to use
Use when an effect should complete successfully with a specific value without any errors
or external dependencies.
Example (Creating a successful effect)
import { Effect } from "effect"
// Creating an effect that represents a successful scenario
//
// ┌─── Effect<number, never, never>
// ▼
const success = Effect.succeed(42)
succeed(import EffectEffect.const onInterrupt: {
<XE, XR>(
finalizer: (
interruptors: ReadonlySet<number>
) => Effect<void, XE, XR>
): <A, E, R>(
self: Effect<A, E, R>
) => Effect<A, E | XE, R | XR>
<A, E, R, XE, XR>(
self: Effect<A, E, R>,
finalizer: (
interruptors: ReadonlySet<number>
) => Effect<void, XE, XR>
): Effect<A, E | XE, R | XR>
}
onInterrupt(import PubSubPubSub.const takeAll: <A>(
self: Subscription<A>
) => Effect.Effect<Arr.NonEmptyArray<A>>
Takes all available messages from the subscription, suspending if no items
are available.
Example (Taking all available messages)
import { Effect, PubSub } from "effect"
const program = Effect.gen(function*() {
const pubsub = yield* PubSub.bounded<string>(10)
yield* Effect.scoped(Effect.gen(function*() {
const subscription = yield* PubSub.subscribe(pubsub)
// Publish multiple messages
yield* PubSub.publishAll(pubsub, ["msg1", "msg2", "msg3"])
// Take all available messages at once
const allMessages = yield* PubSub.takeAll(subscription)
console.log("All messages:", allMessages) // ["msg1", "msg2", "msg3"]
}))
})
takeAll(subscription: PubSub.Subscription<A>(parameter) subscription: {
pubsub: PubSub.Atomic<any>;
subscribers: PubSub.Subscribers<any>;
subscription: PubSub.BackingSubscription<A>;
pollers: MutableList.MutableList<Deferred.Deferred<any>>;
shutdownHook: Latch.Latch;
shutdownFlag: MutableRef.MutableRef<boolean>;
strategy: PubSub.Strategy<any>;
replayWindow: PubSub.ReplayWindow<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; <…;
}
subscription), () => import CauseCause.done())))