Stream
Combinators
mergeWithTag
Signature
declare const mergeWithTag: {
<
S extends {
[key: string]: Stream<any, any, any>;
},
>(
streams: S,
options: {
readonly bufferSize?: number;
readonly concurrency: number | "unbounded";
},
): Stream<
{
[K in string | number | symbol]: {
_tag: K;
value: Stream.Success<S[K]>;
};
}[keyof S],
Error<S[keyof S]>,
Context<S[keyof S]>
>;
(options: { readonly bufferSize?: number; readonly concurrency: number | "unbounded" }): <
S extends {
[key: string]: Stream<any, any, any>;
},
>(
streams: S,
) => Stream<
{
[K in string | number | symbol]: {
_tag: K;
value: Stream.Success<S[K]>;
};
}[keyof S],
Error<S[keyof S]>,
Context<S[keyof S]>
>;
};Example
import { Stream } from "effect"
// Stream.Stream<{ _tag: "a"; value: number; } | { _tag: "b"; value: string; }>
const res = Stream.mergeWithTag(
{
a: Stream.make(0),
b: Stream.make(""),
},
{ concurrency: "unbounded" },
)splitLines
Splits strings on newlines. Handles both Windows newlines (\r\n) and UNIX newlines (\n).
Signature
declare const splitLines: <E, R>(self: Stream<string, E, R>) => Stream<string, E, R>;Constants
DefaultChunkSize
The default chunk size used by the various combinators and constructors of Stream.
Signature
declare const DefaultChunkSize: number;Constructors
acquireRelease
Creates a stream from a single value that will get cleaned up after the stream is consumed.
Signature
declare const acquireRelease: <A, E, R, R2, X>(
acquire: Effect.Effect<A, E, R>,
release: (resource: A, exit: Exit.Exit<unknown, unknown>) => Effect.Effect<X, never, R2>,
) => Stream<A, E, R | R2>;Example
import { Console, Effect, Stream } from "effect"
// Simulating File operations
const open = (filename: string) =>
Effect.gen(function* () {
yield* Console.log(`Opening ${filename}`)
return {
getLines: Effect.succeed(["Line 1", "Line 2", "Line 3"]),
close: Console.log(`Closing ${filename}`),
}
})
const stream = Stream.acquireRelease(open("file.txt"), (file) => file.close).pipe(
Stream.flatMap((file) => file.getLines),
)
Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// Opening file.txt
// Closing file.txt
// { _id: 'Chunk', values: [ [ 'Line 1', 'Line 2', 'Line 3' ] ] }asyncEffect
Creates a stream from an asynchronous callback that can be called multiple times The registration of the callback itself returns an effect. The optionality of the error type E can be used to signal the end of the stream, by setting it to None.
Signature
declare const asyncEffect: <A, E = never, R = never>(
register: (emit: Emit.Emit<R, E, A, void>) => Effect.Effect<unknown, E, R>,
bufferSize?:
| number
| "unbounded"
| {
readonly bufferSize?: number;
readonly strategy?: "dropping" | "sliding" | "suspend";
},
) => Stream<A, E, R>;Creates a stream from an external push-based resource.
You can use the emit helper to emit values to the stream. The emit helper returns a boolean indicating whether the value was emitted or not.
You can also use the emit helper to signal the end of the stream by using apis such as emit.end or emit.fail.
By default it uses an "unbounded" buffer size. You can customize the buffer size and strategy by passing an object as the second argument with the bufferSize and strategy fields.
Signature
declare const asyncPush: <A, E = never, R = never>(
register: (emit: Emit.EmitOpsPush<E, A>) => Effect.Effect<unknown, E, R | Scope.Scope>,
options?:
| {
readonly bufferSize: "unbounded";
}
| {
readonly bufferSize?: number;
readonly strategy?: "dropping" | "sliding";
},
) => Stream<A, E, Exclude<R, Scope.Scope>>;Example
import { Effect, Stream } from "effect"
Stream.asyncPush<string>(
(emit) =>
Effect.acquireRelease(
Effect.gen(function* () {
yield* Effect.log("subscribing")
return setInterval(() => emit.single("tick"), 1000)
}),
(handle) =>
Effect.gen(function* () {
yield* Effect.log("unsubscribing")
clearInterval(handle)
}),
),
{ bufferSize: 16, strategy: "dropping" },
)asyncScoped
Creates a stream from an asynchronous callback that can be called multiple times. The registration of the callback itself returns an a scoped resource. The optionality of the error type E can be used to signal the end of the stream, by setting it to None.
Signature
declare const asyncScoped: <A, E = never, R = never>(
register: (emit: Emit.Emit<R, E, A, void>) => Effect.Effect<unknown, E, R | Scope.Scope>,
bufferSize?:
| number
| "unbounded"
| {
readonly bufferSize?: number;
readonly strategy?: "dropping" | "sliding" | "suspend";
},
) => Stream<A, E, Exclude<R, Scope.Scope>>;Concatenates all of the streams in the chunk to one stream.
Signature
declare const concatAll: <A, E, R>(streams: Chunk.Chunk<Stream<A, E, R>>) => Stream<A, E, R>;Example
import { Chunk, Effect, Stream } from "effect"
const s1 = Stream.make(1, 2, 3)
const s2 = Stream.make(4, 5)
const s3 = Stream.make(6, 7, 8)
const stream = Stream.concatAll(Chunk.make(s1, s2, s3))
Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// {
// _id: 'Chunk',
// values: [
// 1, 2, 3, 4,
// 5, 6, 7, 8
// ]
// }The stream that dies with the specified defect.
Signature
declare const die: (defect: unknown) => Stream<never>;dieMessage
The stream that dies with an exception described by message.
Signature
declare const dieMessage: (message: string) => Stream<never>;The stream that dies with the specified lazily evaluated defect.
Signature
declare const dieSync: (evaluate: LazyArg<unknown>) => Stream<never>;The empty stream.
Signature
declare const empty: Stream<never>;Example
import { Effect, Stream } from "effect"
const stream = Stream.empty
Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// { _id: 'Chunk', values: [] }Creates a stream that executes the specified effect but emits no elements.
Signature
declare const execute: <X, E, R>(effect: Effect.Effect<X, E, R>) => Stream<never, E, R>;Terminates with the specified error.
Signature
declare const fail: <E>(error: E) => Stream<never, E>;Example
import { Effect, Stream } from "effect"
const stream = Stream.fail("Uh oh!")
Effect.runPromiseExit(Stream.runCollect(stream)).then(console.log)
// {
// _id: 'Exit',
// _tag: 'Failure',
// cause: { _id: 'Cause', _tag: 'Fail', failure: 'Uh oh!' }
// }The stream that always fails with the specified Cause.
Signature
declare const failCause: <E>(cause: Cause.Cause<E>) => Stream<never, E>;failCauseSync
The stream that always fails with the specified lazily evaluated Cause.
Signature
declare const failCauseSync: <E>(evaluate: LazyArg<Cause.Cause<E>>) => Stream<never, E>;Terminates with the specified lazily evaluated error.
Signature
declare const failSync: <E>(evaluate: LazyArg<E>) => Stream<never, E>;Creates a one-element stream that never fails and executes the finalizer when it ends.
Signature
declare const finalizer: <R, X>(finalizer: Effect.Effect<X, never, R>) => Stream<void, never, R>;Example
import { Console, Effect, Stream } from "effect"
const application = Stream.fromEffect(Console.log("Application Logic."))
const deleteDir = (dir: string) => Console.log(`Deleting dir: ${dir}`)
const program = application.pipe(
Stream.concat(
Stream.finalizer(
deleteDir("tmp").pipe(Effect.andThen(Console.log("Temporary directory was deleted."))),
),
),
)
Effect.runPromise(Stream.runCollect(program)).then(console.log)
// Application Logic.
// Deleting dir: tmp
// Temporary directory was deleted.
// { _id: 'Chunk', values: [ undefined, undefined ] }fromAsyncIterable
Creates a stream from an AsyncIterable.
Signature
declare const fromAsyncIterable: <A, E>(
iterable: AsyncIterable<A>,
onError: (e: unknown) => E,
) => Stream<A, E>;Example
import { Effect, Stream } from "effect"
const myAsyncIterable = async function* () {
yield 1
yield 2
}
const stream = Stream.fromAsyncIterable(
myAsyncIterable(),
(e) => new Error(String(e)), // Error Handling
)
Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// { _id: 'Chunk', values: [ 1, 2 ] }fromChannel
Creates a stream from a Channel.
Signature
declare const fromChannel: <A, E, R>(
channel: Channel.Channel<Chunk.Chunk<A>, unknown, E, unknown, unknown, unknown, R>,
) => Stream<A, E, R>;Creates a stream from a Chunk of values.
Signature
declare const fromChunk: <A>(chunk: Chunk.Chunk<A>) => Stream<A>;Example
import { Chunk, Effect, Stream } from "effect"
// Creating a stream with values from a single Chunk
const stream = Stream.fromChunk(Chunk.make(1, 2, 3))
Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// { _id: 'Chunk', values: [ 1, 2, 3 ] }fromChunkPubSub
Creates a stream from a subscription to a PubSub.
Options
- shutdown: If true, the PubSub will be shutdown after the stream is evaluated (defaults to false)
Signature
declare const fromChunkPubSub: {
<A>(
pubsub: PubSub<Chunk<A>>,
options: {
readonly scoped: true;
readonly shutdown?: boolean;
},
): Effect<Stream<A, never, never>, never, Scope>;
<A>(
pubsub: PubSub<Chunk<A>>,
options?: {
readonly scoped?: false;
readonly shutdown?: boolean;
},
): Stream<A>;
};fromChunkQueue
Creates a stream from a Queue of values.
Options
- shutdown: If true, the queue will be shutdown after the stream is evaluated (defaults to false)
Signature
declare const fromChunkQueue: <A>(
queue: Queue.Dequeue<Chunk.Chunk<A>>,
options?: {
readonly shutdown?: boolean;
},
) => Stream<A>;fromChunks
Creates a stream from an arbitrary number of chunks.
Signature
declare const fromChunks: <A>(...chunks: Array<Chunk.Chunk<A>>) => Stream<A>;Example
import { Chunk, Effect, Stream } from "effect"
// Creating a stream with values from multiple Chunks
const stream = Stream.fromChunks(Chunk.make(1, 2, 3), Chunk.make(4, 5, 6))
Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// { _id: 'Chunk', values: [ 1, 2, 3, 4, 5, 6 ] }fromEffect
Either emits the success value of this effect or terminates the stream with the failure value of this effect.
Signature
declare const fromEffect: <A, E, R>(effect: Effect.Effect<A, E, R>) => Stream<A, E, R>;Example
import { Effect, Random, Stream } from "effect"
const stream = Stream.fromEffect(Random.nextInt)
Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// Example Output: { _id: 'Chunk', values: [ 922694024 ] }fromEffectOption
Creates a stream from an effect producing a value of type A or an empty Stream.
Signature
declare const fromEffectOption: <A, E, R>(
effect: Effect.Effect<A, Option.Option<E>, R>,
) => Stream<A, E, R>;fromIterable
Creates a new Stream from an iterable collection of values.
Signature
declare const fromIterable: <A>(iterable: Iterable<A>) => Stream<A>;Example
import { Effect, Stream } from "effect"
const numbers = [1, 2, 3]
const stream = Stream.fromIterable(numbers)
Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// { _id: 'Chunk', values: [ 1, 2, 3 ] }fromIterableEffect
Creates a stream from an effect producing a value of type Iterable<A>.
Signature
declare const fromIterableEffect: <A, E, R>(
effect: Effect.Effect<Iterable<A>, E, R>,
) => Stream<A, E, R>;Example
import { Context, Effect, Stream } from "effect"
class Database extends Context.Tag("Database")<
Database,
{ readonly getUsers: Effect.Effect<Array<string>> }
>() {}
const getUsers = Database.pipe(Effect.andThen((_) => _.getUsers))
const stream = Stream.fromIterableEffect(getUsers)
Effect.runPromise(
Stream.runCollect(
stream.pipe(Stream.provideService(Database, { getUsers: Effect.succeed(["user1", "user2"]) })),
),
).then(console.log)
// { _id: 'Chunk', values: [ 'user1', 'user2' ] }fromIteratorSucceed
Creates a stream from an iterator
Signature
declare const fromIteratorSucceed: <A>(
iterator: IterableIterator<A>,
maxChunkSize?: number,
) => Stream<A>;fromPubSub
Creates a stream from a subscription to a PubSub.
Options
- shutdown: If true, the PubSub will be shutdown after the stream is evaluated (defaults to false)
Signature
declare const fromPubSub: {
<A>(
pubsub: PubSub<A>,
options: {
readonly maxChunkSize?: number;
readonly scoped: true;
readonly shutdown?: boolean;
},
): Effect<Stream<A, never, never>, never, Scope>;
<A>(
pubsub: PubSub<A>,
options?: {
readonly maxChunkSize?: number;
readonly scoped?: false;
readonly shutdown?: boolean;
},
): Stream<A>;
};Creates a stream from an effect that pulls elements from another stream.
See Stream.toPull for reference.
Signature
declare const fromPull: <R, R2, E, A>(
effect: Effect.Effect<
Effect.Effect<Chunk.Chunk<A>, Option.Option<E>, R2>,
never,
Scope.Scope | R
>,
) => Stream<A, E, R2 | Exclude<R, Scope.Scope>>;Creates a stream from a queue of values
Options
- maxChunkSize: The maximum number of queued elements to put in one chunk in the stream - shutdown: If true, the queue will be shutdown after the stream is evaluated (defaults to false)
Signature
declare const fromQueue: <A>(
queue: Queue.Dequeue<A>,
options?: {
readonly maxChunkSize?: number;
readonly shutdown?: boolean;
},
) => Stream<A>;fromReadableStream
Creates a stream from a ReadableStream.
See https://developer.mozilla.org/en-US/docs/Web/API/ReadableStream.
Signature
declare const fromReadableStream: {
<A, E>(options: {
readonly evaluate: LazyArg<ReadableStream<A>>;
readonly onError: (error: unknown) => E;
readonly releaseLockOnEnd?: boolean;
}): Stream<A, E>;
<A, E>(evaluate: LazyArg<ReadableStream<A>>, onError: (error: unknown) => E): Stream<A, E>;
};fromReadableStreamByob
Creates a stream from a ReadableStreamBYOBReader.
See https://developer.mozilla.org/en-US/docs/Web/API/ReadableStreamBYOBReader.
Signature
declare const fromReadableStreamByob: {
<E>(options: {
readonly bufferSize?: number;
readonly evaluate: LazyArg<ReadableStream<Uint8Array>>;
readonly onError: (error: unknown) => E;
readonly releaseLockOnEnd?: boolean;
}): Stream<Uint8Array<ArrayBufferLike>, E>;
<E>(
evaluate: LazyArg<ReadableStream<Uint8Array<ArrayBufferLike>>>,
onError: (error: unknown) => E,
allocSize?: number,
): Stream<Uint8Array<ArrayBufferLike>, E>;
};fromSchedule
Creates a stream from a Schedule that does not require any further input. The stream will emit an element for each value output from the schedule, continuing for as long as the schedule continues.
Signature
declare const fromSchedule: <A, R>(
schedule: Schedule.Schedule<A, unknown, R>,
) => Stream<A, never, R>;Example
import { Effect, Schedule, Stream } from "effect"
// Emits values every 1 second for a total of 5 emissions
const schedule = Schedule.spaced("1 second").pipe(Schedule.compose(Schedule.recurs(5)))
const stream = Stream.fromSchedule(schedule)
Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// { _id: 'Chunk', values: [ 0, 1, 2, 3, 4 ] }fromTPubSub
Creates a stream from a subscription to a TPubSub.
Signature
declare const fromTPubSub: <A>(pubsub: TPubSub<A>) => Stream<A>;fromTQueue
Creates a stream from a TQueue of values
Signature
declare const fromTQueue: <A>(queue: TDequeue<A>) => Stream<A>;The infinite stream of iterative function application: a, f(a), f(f(a)), f(f(f(a))), ...
Signature
declare const iterate: <A>(value: A, next: (value: A) => A) => Stream<A>;Example
import { Effect, Stream } from "effect"
// An infinite Stream of numbers starting from 1 and incrementing
const stream = Stream.iterate(1, (n) => n + 1)
Effect.runPromise(Stream.runCollect(stream.pipe(Stream.take(10)))).then(console.log)
// { _id: 'Chunk', values: [ 1, 2, 3, 4, 5, 6, 7, 8, 9, 10 ] }Creates a stream from an sequence of values.
Signature
declare const make: <As extends Array<any>>(...as: As) => Stream<As[number]>;Example
import { Effect, Stream } from "effect"
const stream = Stream.make(1, 2, 3)
Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// { _id: 'Chunk', values: [ 1, 2, 3 ] }The stream that never produces any value or fails with any error.
Signature
declare const never: Stream<never>;Like Stream.unfold, but allows the emission of values to end one step further than the unfolding of the state. This is useful for embedding paginated APIs, hence the name.
Signature
declare const paginate: <S, A>(s: S, f: (s: S) => readonly [A, Option.Option<S>]) => Stream<A>;Example
import { Effect, Option, Stream } from "effect"
const stream = Stream.paginate(0, (n) => [n, n < 3 ? Option.some(n + 1) : Option.none()])
Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// { _id: 'Chunk', values: [ 0, 1, 2, 3 ] }paginateChunk
Like Stream.unfoldChunk, but allows the emission of values to end one step further than the unfolding of the state. This is useful for embedding paginated APIs, hence the name.
Signature
declare const paginateChunk: <S, A>(
s: S,
f: (s: S) => readonly [Chunk.Chunk<A>, Option.Option<S>],
) => Stream<A>;paginateChunkEffect
Like Stream.unfoldChunkEffect, but allows the emission of values to end one step further than the unfolding of the state. This is useful for embedding paginated APIs, hence the name.
Signature
declare const paginateChunkEffect: <S, A, E, R>(
s: S,
f: (s: S) => Effect.Effect<readonly [Chunk.Chunk<A>, Option.Option<S>], E, R>,
) => Stream<A, E, R>;paginateEffect
Like Stream.unfoldEffect but allows the emission of values to end one step further than the unfolding of the state. This is useful for embedding paginated APIs, hence the name.
Signature
declare const paginateEffect: <S, A, E, R>(
s: S,
f: (s: S) => Effect.Effect<readonly [A, Option.Option<S>], E, R>,
) => Stream<A, E, R>;Constructs a stream from a range of integers, including both endpoints.
Signature
declare const range: (min: number, max: number, chunkSize?: number) => Stream<number>;Example
import { Effect, Stream } from "effect"
// A Stream with a range of numbers from 1 to 5
const stream = Stream.range(1, 5)
Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// { _id: 'Chunk', values: [ 1, 2, 3, 4, 5 ] }repeatEffect
Creates a stream from an effect producing a value of type A which repeats forever.
Signature
declare const repeatEffect: <A, E, R>(effect: Effect.Effect<A, E, R>) => Stream<A, E, R>;Example
import { Effect, Random, Stream } from "effect"
const stream = Stream.repeatEffect(Random.nextInt)
Effect.runPromise(Stream.runCollect(stream.pipe(Stream.take(5)))).then(console.log)
// Example Output: { _id: 'Chunk', values: [ 3891571149, 4239494205, 2352981603, 2339111046, 1488052210 ] }repeatEffectChunk
Creates a stream from an effect producing chunks of A values which repeats forever.
Signature
declare const repeatEffectChunk: <A, E, R>(
effect: Effect.Effect<Chunk.Chunk<A>, E, R>,
) => Stream<A, E, R>;repeatEffectChunkOption
Creates a stream from an effect producing chunks of A values until it fails with None.
Signature
declare const repeatEffectChunkOption: <A, E, R>(
effect: Effect.Effect<Chunk.Chunk<A>, Option.Option<E>, R>,
) => Stream<A, E, R>;repeatEffectOption
Creates a stream from an effect producing values of type A until it fails with None.
Signature
declare const repeatEffectOption: <A, E, R>(
effect: Effect.Effect<A, Option.Option<E>, R>,
) => Stream<A, E, R>;Example
// In this example, we're draining an Iterator to create a stream from it
import { Stream, Effect, Option } from "effect"
const drainIterator = <A>(it: Iterator<A>): Stream.Stream<A> =>
Stream.repeatEffectOption(
Effect.sync(() => it.next()).pipe(
Effect.andThen((res) => {
if (res.done) {
return Effect.fail(Option.none())
}
return Effect.succeed(res.value)
}),
),
)repeatEffectWithSchedule
Creates a stream from an effect producing a value of type A, which is repeated using the specified schedule.
Signature
declare const repeatEffectWithSchedule: <A, E, R, X, A0 extends A, R2>(
effect: Effect.Effect<A, E, R>,
schedule: Schedule.Schedule<X, A0, R2>,
) => Stream<A, E, R | R2>;repeatValue
Repeats the provided value infinitely.
Signature
declare const repeatValue: <A>(value: A) => Stream<A>;Example
import { Effect, Stream } from "effect"
const stream = Stream.repeatValue(0)
Effect.runPromise(Stream.runCollect(stream.pipe(Stream.take(5)))).then(console.log)
// { _id: 'Chunk', values: [ 0, 0, 0, 0, 0 ] }Creates a single-valued stream from a scoped resource.
Signature
declare const scoped: <A, E, R>(
effect: Effect.Effect<A, E, R>,
) => Stream<A, E, Exclude<R, Scope.Scope>>;Example
import { Console, Effect, Stream } from "effect"
// Creating a single-valued stream from a scoped resource
const stream = Stream.scoped(
Effect.acquireRelease(Console.log("acquire"), () => Console.log("release")),
).pipe(Stream.flatMap(() => Console.log("use")))
Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// acquire
// use
// release
// { _id: 'Chunk', values: [ undefined ] }scopedWith
Use a function that receives a scope and returns an effect to emit an output element. The output element will be the result of the returned effect, if successful.
Signature
declare const scopedWith: <A, E, R>(
f: (scope: Scope.Scope) => Effect.Effect<A, E, R>,
) => Stream<A, E, R>;Creates a single-valued pure stream.
Signature
declare const succeed: <A>(value: A) => Stream<A>;Example
import { Effect, Stream } from "effect"
// A Stream with a single number
const stream = Stream.succeed(3)
Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// { _id: 'Chunk', values: [ 3 ] }Returns a lazily constructed stream.
Signature
declare const suspend: <A, E, R>(stream: LazyArg<Stream<A, E, R>>) => Stream<A, E, R>;Creates a single-valued pure stream.
Signature
declare const sync: <A>(evaluate: LazyArg<A>) => Stream<A>;A stream that emits void values spaced by the specified duration.
Signature
declare const tick: (interval: Duration.DurationInput) => Stream<void>;Example
import { Effect, Stream } from "effect"
let last = Date.now()
const log = (message: string) =>
Effect.sync(() => {
const end = Date.now()
console.log(`${message} after ${end - last}ms`)
last = end
})
const stream = Stream.tick("1 seconds").pipe(Stream.tap(() => log("tick")))
Effect.runPromise(Stream.runCollect(stream.pipe(Stream.take(5)))).then(console.log)
// tick after 4ms
// tick after 1003ms
// tick after 1001ms
// tick after 1002ms
// tick after 1002ms
// { _id: 'Chunk', values: [ undefined, undefined, undefined, undefined, undefined ] }Creates a channel from a Stream.
Signature
declare const toChannel: <A, E, R>(
stream: Stream<A, E, R>,
) => Channel.Channel<Chunk.Chunk<A>, unknown, E, unknown, unknown, unknown, R>;Creates a stream by peeling off the "layers" of a value of type S.
Signature
declare const unfold: <S, A>(s: S, f: (s: S) => Option.Option<readonly [A, S]>) => Stream<A>;Example
import { Effect, Option, Stream } from "effect"
const stream = Stream.unfold(1, (n) => Option.some([n, n + 1]))
Effect.runPromise(Stream.runCollect(stream.pipe(Stream.take(5)))).then(console.log)
// { _id: 'Chunk', values: [ 1, 2, 3, 4, 5 ] }unfoldChunk
Creates a stream by peeling off the "layers" of a value of type S.
Signature
declare const unfoldChunk: <S, A>(
s: S,
f: (s: S) => Option.Option<readonly [Chunk.Chunk<A>, S]>,
) => Stream<A>;unfoldChunkEffect
Creates a stream by effectfully peeling off the "layers" of a value of type S.
Signature
declare const unfoldChunkEffect: <S, A, E, R>(
s: S,
f: (s: S) => Effect.Effect<Option.Option<readonly [Chunk.Chunk<A>, S]>, E, R>,
) => Stream<A, E, R>;unfoldEffect
Creates a stream by effectfully peeling off the "layers" of a value of type S.
Signature
declare const unfoldEffect: <S, A, E, R>(
s: S,
f: (s: S) => Effect.Effect<Option.Option<readonly [A, S]>, E, R>,
) => Stream<A, E, R>;Example
import { Effect, Option, Random, Stream } from "effect"
const stream = Stream.unfoldEffect(1, (n) =>
Random.nextBoolean.pipe(Effect.map((b) => (b ? Option.some([n, -n]) : Option.some([n, n])))),
)
Effect.runPromise(Stream.runCollect(stream.pipe(Stream.take(5)))).then(console.log)
// { _id: 'Chunk', values: [ 1, -1, -1, -1, -1 ] }Creates a stream produced from an Effect.
Signature
declare const unwrap: <A, E2, R2, E, R>(
effect: Effect.Effect<Stream<A, E2, R2>, E, R>,
) => Stream<A, E | E2, R | R2>;unwrapScoped
Creates a stream produced from a scoped Effect.
Signature
declare const unwrapScoped: <A, E2, R2, E, R>(
effect: Effect.Effect<Stream<A, E2, R2>, E, R>,
) => Stream<A, E | E2, R2 | Exclude<R, Scope.Scope>>;unwrapScopedWith
Creates a stream produced from a function which receives a Scope and returns an Effect. The resulting stream will emit a single element, which will be the result of the returned effect, if successful.
Signature
declare const unwrapScopedWith: <A, E2, R2, E, R>(
f: (scope: Scope.Scope) => Effect.Effect<Stream<A, E2, R2>, E, R>,
) => Stream<A, E | E2, R | R2>;Returns the resulting stream when the given PartialFunction is defined for the given value, otherwise returns an empty stream.
Signature
declare const whenCase: <A, A2, E, R>(
evaluate: LazyArg<A>,
pf: (a: A) => Option.Option<Stream<A2, E, R>>,
) => Stream<A2, E, R>;Context
Accesses the whole context of the stream.
Signature
declare const context: <R>() => Stream<Context.Context<R>, never, R>;contextWith
Accesses the context of the stream.
Signature
declare const contextWith: <R, A>(f: (env: Context.Context<R>) => A) => Stream<A, never, R>;contextWithEffect
Accesses the context of the stream in the context of an effect.
Signature
declare const contextWithEffect: <R0, A, E, R>(
f: (env: Context.Context<R0>) => Effect.Effect<A, E, R>,
) => Stream<A, E, R0 | R>;contextWithStream
Accesses the context of the stream in the context of a stream.
Signature
declare const contextWithStream: <R0, A, E, R>(
f: (env: Context.Context<R0>) => Stream<A, E, R>,
) => Stream<A, E, R0 | R>;mapInputContext
Transforms the context being provided to the stream with the specified function.
Signature
declare const mapInputContext: {
<R0, R>(f: (env: Context<R0>) => Context<R>): <A, E>(self: Stream<A, E, R>) => Stream<A, E, R0>;
<A, E, R0, R>(self: Stream<A, E, R>, f: (env: Context<R0>) => Context<R>): Stream<A, E, R0>;
};provideContext
Provides the stream with its required context, which eliminates its dependency on R.
Signature
declare const provideContext: {
<R>(context: Context<R>): <A, E>(self: Stream<A, E, R>) => Stream<A, E>;
<A, E, R>(self: Stream<A, E, R>, context: Context<R>): Stream<A, E>;
};provideLayer
Provides a Layer to the stream, which translates it to another level.
Signature
declare const provideLayer: {
<RIn, E2, ROut>(
layer: Layer<ROut, E2, RIn>,
): <A, E>(self: Stream<A, E, ROut>) => Stream<A, E2 | E, RIn>;
<A, E, RIn, E2, ROut>(
self: Stream<A, E, ROut>,
layer: Layer<ROut, E2, RIn>,
): Stream<A, E | E2, RIn>;
};provideService
Provides the stream with the single service it requires. If the stream requires more than one service use Stream.provideContext instead.
Signature
declare const provideService: {
<I, S>(
tag: Tag<I, S>,
resource: NoInfer<S>,
): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E, Exclude<R, I>>;
<A, E, R, I, S>(
self: Stream<A, E, R>,
tag: Tag<I, S>,
resource: NoInfer<S>,
): Stream<A, E, Exclude<R, I>>;
};provideServiceEffect
Provides the stream with the single service it requires. If the stream requires more than one service use Stream.provideContext instead.
Signature
declare const provideServiceEffect: {
<I, S, E2, R2>(
tag: Tag<I, S>,
effect: Effect<NoInfer<S>, E2, R2>,
): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E2 | E, R2 | Exclude<R, I>>;
<A, E, R, I, S, E2, R2>(
self: Stream<A, E, R>,
tag: Tag<I, S>,
effect: Effect<NoInfer<S>, E2, R2>,
): Stream<A, E | E2, R2 | Exclude<R, I>>;
};provideServiceStream
Provides the stream with the single service it requires. If the stream requires more than one service use Stream.provideContext instead.
Signature
declare const provideServiceStream: {
<I, S, E2, R2>(
tag: Tag<I, S>,
stream: Stream<NoInfer<S>, E2, R2>,
): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E2 | E, R2 | Exclude<R, I>>;
<A, E, R, I, S, E2, R2>(
self: Stream<A, E, R>,
tag: Tag<I, S>,
stream: Stream<NoInfer<S>, E2, R2>,
): Stream<A, E | E2, R2 | Exclude<R, I>>;
};provideSomeContext
Provides the stream with some of its required context, which eliminates its dependency on R.
Signature
declare const provideSomeContext: {
<R2>(context: Context<R2>): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E, Exclude<R, R2>>;
<A, E, R, R2>(self: Stream<A, E, R>, context: Context<R2>): Stream<A, E, Exclude<R, R2>>;
};provideSomeLayer
Splits the context into two parts, providing one part using the specified layer and leaving the remainder R0.
Signature
declare const provideSomeLayer: {
<RIn, E2, ROut>(
layer: Layer<ROut, E2, RIn>,
): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E2 | E, RIn | Exclude<R, ROut>>;
<A, E, R, RIn, E2, ROut>(
self: Stream<A, E, R>,
layer: Layer<ROut, E2, RIn>,
): Stream<A, E | E2, RIn | Exclude<R, ROut>>;
};updateService
Updates the specified service within the context of the Stream.
Signature
declare const updateService: {
<I, S>(
tag: Tag<I, S>,
f: (service: NoInfer<S>) => NoInfer<S>,
): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E, I | R>;
<A, E, R, I, S>(
self: Stream<A, E, R>,
tag: Tag<I, S>,
f: (service: NoInfer<S>) => NoInfer<S>,
): Stream<A, E, R | I>;
};Destructors
Runs the sink on the stream to produce either the sink's result or an error.
Signature
declare const run: {
<A2, A, E2, R2>(
sink: Sink<A2, A, unknown, E2, R2>,
): <E, R>(self: Stream<A, E, R>) => Effect<A2, E2 | E, Exclude<R2, Scope> | Exclude<R, Scope>>;
<A, E, R, A2, E2, R2>(
self: Stream<A, E, R>,
sink: Sink<A2, A, unknown, E2, R2>,
): Effect<A2, E | E2, Exclude<R, Scope> | Exclude<R2, Scope>>;
};runCollect
Runs the stream and collects all of its elements to a chunk.
Signature
declare const runCollect: <A, E, R>(self: Stream<A, E, R>) => Effect.Effect<Chunk.Chunk<A>, E, R>;Runs the stream and emits the number of elements processed
Signature
declare const runCount: <A, E, R>(self: Stream<A, E, R>) => Effect.Effect<number, E, R>;Runs the stream only for its effects. The emitted elements are discarded.
Signature
declare const runDrain: <A, E, R>(self: Stream<A, E, R>) => Effect.Effect<void, E, R>;Executes a pure fold over the stream of values - reduces all elements in the stream to a value of type S.
Signature
declare const runFold: {
<S, A>(s: S, f: (s: S, a: A) => S): <E, R>(self: Stream<A, E, R>) => Effect<S, E, R>;
<A, E, R, S>(self: Stream<A, E, R>, s: S, f: (s: S, a: A) => S): Effect<S, E, R>;
};runFoldEffect
Executes an effectful fold over the stream of values.
Signature
declare const runFoldEffect: {
<S, A, E2, R2>(
s: S,
f: (s: S, a: A) => Effect<S, E2, R2>,
): <E, R>(self: Stream<A, E, R>) => Effect<S, E2 | E, Exclude<R2, Scope> | Exclude<R, Scope>>;
<A, E, R, S, E2, R2>(
self: Stream<A, E, R>,
s: S,
f: (s: S, a: A) => Effect<S, E2, R2>,
): Effect<S, E | E2, Exclude<R, Scope> | Exclude<R2, Scope>>;
};runFoldScoped
Executes a pure fold over the stream of values. Returns a scoped value that represents the scope of the stream.
Signature
declare const runFoldScoped: {
<S, A>(s: S, f: (s: S, a: A) => S): <E, R>(self: Stream<A, E, R>) => Effect<S, E, Scope | R>;
<A, E, R, S>(self: Stream<A, E, R>, s: S, f: (s: S, a: A) => S): Effect<S, E, Scope | R>;
};runFoldScopedEffect
Executes an effectful fold over the stream of values. Returns a scoped value that represents the scope of the stream.
Signature
declare const runFoldScopedEffect: {
<S, A, E2, R2>(
s: S,
f: (s: S, a: A) => Effect<S, E2, R2>,
): <E, R>(self: Stream<A, E, R>) => Effect<S, E2 | E, Scope | R2 | R>;
<A, E, R, S, E2, R2>(
self: Stream<A, E, R>,
s: S,
f: (s: S, a: A) => Effect<S, E2, R2>,
): Effect<S, E | E2, Scope | R | R2>;
};runFoldWhile
Reduces the elements in the stream to a value of type S. Stops the fold early when the condition is not fulfilled. Example:
Signature
declare const runFoldWhile: {
<S, A>(
s: S,
cont: Predicate<S>,
f: (s: S, a: A) => S,
): <E, R>(self: Stream<A, E, R>) => Effect<S, E, R>;
<A, E, R, S>(
self: Stream<A, E, R>,
s: S,
cont: Predicate<S>,
f: (s: S, a: A) => S,
): Effect<S, E, R>;
};runFoldWhileEffect
Executes an effectful fold over the stream of values. Stops the fold early when the condition is not fulfilled.
Signature
declare const runFoldWhileEffect: {
<S, A, E2, R2>(
s: S,
cont: Predicate<S>,
f: (s: S, a: A) => Effect<S, E2, R2>,
): <E, R>(self: Stream<A, E, R>) => Effect<S, E2 | E, Exclude<R2, Scope> | Exclude<R, Scope>>;
<A, E, R, S, E2, R2>(
self: Stream<A, E, R>,
s: S,
cont: Predicate<S>,
f: (s: S, a: A) => Effect<S, E2, R2>,
): Effect<S, E | E2, Exclude<R, Scope> | Exclude<R2, Scope>>;
};runFoldWhileScoped
Executes a pure fold over the stream of values. Returns a scoped value that represents the scope of the stream. Stops the fold early when the condition is not fulfilled.
Signature
declare const runFoldWhileScoped: {
<S, A>(
s: S,
cont: Predicate<S>,
f: (s: S, a: A) => S,
): <E, R>(self: Stream<A, E, R>) => Effect<S, E, Scope | R>;
<A, E, R, S>(
self: Stream<A, E, R>,
s: S,
cont: Predicate<S>,
f: (s: S, a: A) => S,
): Effect<S, E, Scope | R>;
};runFoldWhileScopedEffect
Executes an effectful fold over the stream of values. Returns a scoped value that represents the scope of the stream. Stops the fold early when the condition is not fulfilled.
Signature
declare const runFoldWhileScopedEffect: {
<S, A, E2, R2>(
s: S,
cont: Predicate<S>,
f: (s: S, a: A) => Effect<S, E2, R2>,
): <E, R>(self: Stream<A, E, R>) => Effect<S, E2 | E, Scope | R2 | R>;
<A, E, R, S, E2, R2>(
self: Stream<A, E, R>,
s: S,
cont: Predicate<S>,
f: (s: S, a: A) => Effect<S, E2, R2>,
): Effect<S, E | E2, Scope | R | R2>;
};runForEach
Consumes all elements of the stream, passing them to the specified callback.
Signature
declare const runForEach: {
<A, X, E2, R2>(
f: (a: A) => Effect<X, E2, R2>,
): <E, R>(self: Stream<A, E, R>) => Effect<void, E2 | E, R2 | R>;
<A, E, R, X, E2, R2>(
self: Stream<A, E, R>,
f: (a: A) => Effect<X, E2, R2>,
): Effect<void, E | E2, R | R2>;
};runForEachChunk
Consumes all elements of the stream, passing them to the specified callback.
Signature
declare const runForEachChunk: {
<A, X, E2, R2>(
f: (a: Chunk<A>) => Effect<X, E2, R2>,
): <E, R>(self: Stream<A, E, R>) => Effect<void, E2 | E, R2 | R>;
<A, E, R, X, E2, R2>(
self: Stream<A, E, R>,
f: (a: Chunk<A>) => Effect<X, E2, R2>,
): Effect<void, E | E2, R | R2>;
};runForEachChunkScoped
Like Stream.runForEachChunk, but returns a scoped effect so the finalization order can be controlled.
Signature
declare const runForEachChunkScoped: {
<A, X, E2, R2>(
f: (a: Chunk<A>) => Effect<X, E2, R2>,
): <E, R>(self: Stream<A, E, R>) => Effect<void, E2 | E, Scope | R2 | R>;
<A, E, R, X, E2, R2>(
self: Stream<A, E, R>,
f: (a: Chunk<A>) => Effect<X, E2, R2>,
): Effect<void, E | E2, Scope | R | R2>;
};runForEachScoped
Like Stream.forEach, but returns a scoped effect so the finalization order can be controlled.
Signature
declare const runForEachScoped: {
<A, X, E2, R2>(
f: (a: A) => Effect<X, E2, R2>,
): <E, R>(self: Stream<A, E, R>) => Effect<void, E2 | E, Scope | R2 | R>;
<A, E, R, X, E2, R2>(
self: Stream<A, E, R>,
f: (a: A) => Effect<X, E2, R2>,
): Effect<void, E | E2, Scope | R | R2>;
};runForEachWhile
Consumes elements of the stream, passing them to the specified callback, and terminating consumption when the callback returns false.
Signature
declare const runForEachWhile: {
<A, E2, R2>(
f: (a: A) => Effect<boolean, E2, R2>,
): <E, R>(self: Stream<A, E, R>) => Effect<void, E2 | E, R2 | R>;
<A, E, R, E2, R2>(
self: Stream<A, E, R>,
f: (a: A) => Effect<boolean, E2, R2>,
): Effect<void, E | E2, R | R2>;
};runForEachWhileScoped
Like Stream.runForEachWhile, but returns a scoped effect so the finalization order can be controlled.
Signature
declare const runForEachWhileScoped: {
<A, E2, R2>(
f: (a: A) => Effect<boolean, E2, R2>,
): <E, R>(self: Stream<A, E, R>) => Effect<void, E2 | E, Scope | R2 | R>;
<A, E, R, E2, R2>(
self: Stream<A, E, R>,
f: (a: A) => Effect<boolean, E2, R2>,
): Effect<void, E | E2, Scope | R | R2>;
};Runs the stream to completion and yields the first value emitted by it, discarding the rest of the elements.
Signature
declare const runHead: <A, E, R>(self: Stream<A, E, R>) => Effect.Effect<Option.Option<A>, E, R>;runIntoPubSub
Publishes elements of this stream to a PubSub. Stream failure and ending will also be signalled.
Signature
declare const runIntoPubSub: {
<A, E>(pubsub: PubSub<Take<A, E>>): <R>(self: Stream<A, E, R>) => Effect<void, never, R>;
<A, E, R>(self: Stream<A, E, R>, pubsub: PubSub<Take<A, E>>): Effect<void, never, R>;
};runIntoPubSubScoped
Like Stream.runIntoPubSub, but provides the result as a scoped effect to allow for scope composition.
Signature
declare const runIntoPubSubScoped: {
<A, E>(pubsub: PubSub<Take<A, E>>): <R>(self: Stream<A, E, R>) => Effect<void, never, Scope | R>;
<A, E, R>(self: Stream<A, E, R>, pubsub: PubSub<Take<A, E>>): Effect<void, never, Scope | R>;
};runIntoQueue
Enqueues elements of this stream into a queue. Stream failure and ending will also be signalled.
Signature
declare const runIntoQueue: {
<A, E>(queue: Enqueue<Take<A, E>>): <R>(self: Stream<A, E, R>) => Effect<void, never, R>;
<A, E, R>(self: Stream<A, E, R>, queue: Enqueue<Take<A, E>>): Effect<void, never, R>;
};runIntoQueueElementsScoped
Like Stream.runIntoQueue, but provides the result as a scoped Effect to allow for scope composition.
Signature
declare const runIntoQueueElementsScoped: {
<A, E>(
queue: Enqueue<Exit<A, Option<E>>>,
): <R>(self: Stream<A, E, R>) => Effect<void, never, Scope | R>;
<A, E, R>(
self: Stream<A, E, R>,
queue: Enqueue<Exit<A, Option<E>>>,
): Effect<void, never, Scope | R>;
};runIntoQueueScoped
Like Stream.runIntoQueue, but provides the result as a scoped effect to allow for scope composition.
Signature
declare const runIntoQueueScoped: {
<A, E>(queue: Enqueue<Take<A, E>>): <R>(self: Stream<A, E, R>) => Effect<void, never, Scope | R>;
<A, E, R>(self: Stream<A, E, R>, queue: Enqueue<Take<A, E>>): Effect<void, never, Scope | R>;
};Runs the stream to completion and yields the last value emitted by it, discarding the rest of the elements.
Signature
declare const runLast: <A, E, R>(self: Stream<A, E, R>) => Effect.Effect<Option.Option<A>, E, R>;Signature
declare const runScoped: {
<A2, A, E2, R2>(
sink: Sink<A2, A, unknown, E2, R2>,
): <E, R>(self: Stream<A, E, R>) => Effect<A2, E2 | E, Scope | R2 | R>;
<A, E, R, A2, E2, R2>(
self: Stream<A, E, R>,
sink: Sink<A2, A, unknown, E2, R2>,
): Effect<A2, E | E2, Scope | R | R2>;
};Runs the stream to a sink which sums elements, provided they are Numeric.
Signature
declare const runSum: <E, R>(self: Stream<number, E, R>) => Effect.Effect<number, E, R>;toAsyncIterable
Converts the stream to a AsyncIterable.
Signature
declare const toAsyncIterable: <A, E>(self: Stream<A, E>) => AsyncIterable<A>;toAsyncIterableEffect
Converts the stream to a AsyncIterable capturing the required dependencies.
Signature
declare const toAsyncIterableEffect: <A, E, R>(
self: Stream<A, E, R>,
) => Effect.Effect<AsyncIterable<A>, never, R>;toAsyncIterableRuntime
Converts the stream to a AsyncIterable using the provided runtime.
Signature
declare const toAsyncIterableRuntime: {
<A, XR>(runtime: Runtime<XR>): <E, R>(self: Stream<A, E, R>) => AsyncIterable<A>;
<A, E, XR, R>(self: Stream<A, E, R>, runtime: Runtime<XR>): AsyncIterable<A>;
};Converts the stream to a scoped PubSub of chunks. After the scope is closed, the PubSub will never again produce values and should be discarded.
Signature
declare const toPubSub: {
(
capacity:
| number
| {
readonly capacity: "unbounded";
readonly replay?: number;
}
| {
readonly capacity: number;
readonly replay?: number;
readonly strategy?: "sliding" | "dropping" | "suspend";
},
): <A, E, R>(self: Stream<A, E, R>) => Effect<PubSub<Take<A, E>>, never, Scope | R>;
<A, E, R>(
self: Stream<A, E, R>,
capacity:
| number
| {
readonly capacity: "unbounded";
readonly replay?: number;
}
| {
readonly capacity: number;
readonly replay?: number;
readonly strategy?: "sliding" | "dropping" | "suspend";
},
): Effect<PubSub<Take<A, E>>, never, Scope | R>;
};Returns in a scope an Effect that can be used to repeatedly pull chunks from the stream. The pull effect fails with None when the stream is finished, or with Some error if it fails, otherwise it returns a chunk of the stream's output.
Signature
declare const toPull: <A, E, R>(
self: Stream<A, E, R>,
) => Effect.Effect<Effect.Effect<Chunk.Chunk<A>, Option.Option<E>, R>, never, Scope.Scope | R>;Example
import { Effect, Stream } from "effect"
// Simulate a chunked stream
const stream = Stream.fromIterable([1, 2, 3, 4, 5]).pipe(Stream.rechunk(2))
const program = Effect.gen(function* () {
// Create an effect to get data chunks from the stream
const getChunk = yield* Stream.toPull(stream)
// Continuously fetch and process chunks
while (true) {
const chunk = yield* getChunk
console.log(chunk)
}
})
Effect.runPromise(Effect.scoped(program)).then(console.log, console.error)
// { _id: 'Chunk', values: [ 1, 2 ] }
// { _id: 'Chunk', values: [ 3, 4 ] }
// { _id: 'Chunk', values: [ 5 ] }
// (FiberFailure) Error: {
// "_id": "Option",
// "_tag": "None"
// }Converts the stream to a scoped queue of chunks. After the scope is closed, the queue will never again produce values and should be discarded.
Defaults to the "suspend" back pressure strategy with a capacity of 2.
Signature
declare const toQueue: {
(
options?:
| {
readonly capacity?: number;
readonly strategy?: "dropping" | "sliding" | "suspend";
}
| {
readonly strategy: "unbounded";
},
): <A, E, R>(self: Stream<A, E, R>) => Effect<Dequeue<Take<A, E>>, never, Scope | R>;
<A, E, R>(
self: Stream<A, E, R>,
options?:
| {
readonly capacity?: number;
readonly strategy?: "dropping" | "sliding" | "suspend";
}
| {
readonly strategy: "unbounded";
},
): Effect<Dequeue<Take<A, E>>, never, Scope | R>;
};toQueueOfElements
Converts the stream to a scoped queue of elements. After the scope is closed, the queue will never again produce values and should be discarded.
Defaults to a capacity of 2.
Signature
declare const toQueueOfElements: {
(options?: {
readonly capacity?: number;
}): <A, E, R>(self: Stream<A, E, R>) => Effect<Dequeue<Exit<A, Option<E>>>, never, Scope | R>;
<A, E, R>(
self: Stream<A, E, R>,
options?: {
readonly capacity?: number;
},
): Effect<Dequeue<Exit<A, Option<E>>>, never, Scope | R>;
};toReadableStream
Converts the stream to a ReadableStream.
See https://developer.mozilla.org/en-US/docs/Web/API/ReadableStream.
Signature
declare const toReadableStream: {
<A>(options?: {
readonly strategy?: QueuingStrategy<A>;
}): <E>(self: Stream<A, E>) => ReadableStream<A>;
<A, E>(
self: Stream<A, E>,
options?: {
readonly strategy?: QueuingStrategy<A>;
},
): ReadableStream<A>;
};toReadableStreamEffect
Converts the stream to a Effect<ReadableStream>.
See https://developer.mozilla.org/en-US/docs/Web/API/ReadableStream.
Signature
declare const toReadableStreamEffect: {
<A>(options?: {
readonly strategy?: QueuingStrategy<A>;
}): <E, R>(self: Stream<A, E, R>) => Effect<ReadableStream<A>, never, R>;
<A, E, R>(
self: Stream<A, E, R>,
options?: {
readonly strategy?: QueuingStrategy<A>;
},
): Effect<ReadableStream<A>, never, R>;
};toReadableStreamRuntime
Converts the stream to a ReadableStream using the provided runtime.
See https://developer.mozilla.org/en-US/docs/Web/API/ReadableStream.
Signature
declare const toReadableStreamRuntime: {
<A, XR>(
runtime: Runtime<XR>,
options?: {
readonly strategy?: QueuingStrategy<A>;
},
): <E, R>(self: Stream<A, E, R>) => ReadableStream<A>;
<A, E, XR, R>(
self: Stream<A, E, R>,
runtime: Runtime<XR>,
options?: {
readonly strategy?: QueuingStrategy<A>;
},
): ReadableStream<A>;
};Do Notation
The "do simulation" in Effect allows you to write code in a more declarative style, similar to the "do notation" in other programming languages. It provides a way to define variables and perform operations on them using functions like bind and let.
Here's how the do simulation works:
1. Start the do simulation using the Do value 2. Within the do simulation scope, you can use the bind function to define variables and bind them to Stream values 3. You can accumulate multiple bind statements to define multiple variables within the scope 4. Inside the do simulation scope, you can also use the let function to define variables and bind them to simple values
See
Signature
declare const bind: {
<N extends string, A, B, E2, R2>(
tag: Exclude<N, keyof A>,
f: (_: NoInfer<A>) => Stream<B, E2, R2>,
options?: {
readonly bufferSize?: number;
readonly concurrency?: number | "unbounded";
},
): <E, R>(
self: Stream<A, E, R>,
) => Stream<{ [K in string | number | symbol]: K extends keyof A ? A[K] : B }, E2 | E, R2 | R>;
<A, E, R, N extends string, B, E2, R2>(
self: Stream<A, E, R>,
tag: Exclude<N, keyof A>,
f: (_: NoInfer<A>) => Stream<B, E2, R2>,
options?: {
readonly bufferSize?: number;
readonly concurrency?: number | "unbounded";
},
): Stream<{ [K in string | number | symbol]: K extends keyof A ? A[K] : B }, E | E2, R | R2>;
};Example
import * as assert from "node:assert"
import { Chunk, Effect, pipe, Stream } from "effect"
const result = pipe(
Stream.Do,
Stream.bind("x", () => Stream.succeed(2)),
Stream.bind("y", () => Stream.succeed(3)),
Stream.let("sum", ({ x, y }) => x + y),
)
assert.deepStrictEqual(Effect.runSync(Stream.runCollect(result)), Chunk.of({ x: 2, y: 3, sum: 5 }))bindEffect
Signature
declare const bindEffect: {
<N extends string, A, B, E2, R2>(
tag: Exclude<N, keyof A>,
f: (_: NoInfer<A>) => Effect<B, E2, R2>,
options?: {
readonly bufferSize?: number;
readonly concurrency?: number | "unbounded";
},
): <E, R>(
self: Stream<A, E, R>,
) => Stream<{ [K in string | number | symbol]: K extends keyof A ? A[K] : B }, E2 | E, R2 | R>;
<A, E, R, N extends string, B, E2, R2>(
self: Stream<A, E, R>,
tag: Exclude<N, keyof A>,
f: (_: NoInfer<A>) => Effect<B, E2, R2>,
options?: {
readonly concurrency?: number | "unbounded";
readonly unordered?: boolean;
},
): Stream<{ [K in string | number | symbol]: K extends keyof A ? A[K] : B }, E | E2, R | R2>;
};The "do simulation" in Effect allows you to write code in a more declarative style, similar to the "do notation" in other programming languages. It provides a way to define variables and perform operations on them using functions like bind and let.
Here's how the do simulation works:
1. Start the do simulation using the Do value 2. Within the do simulation scope, you can use the bind function to define variables and bind them to Stream values 3. You can accumulate multiple bind statements to define multiple variables within the scope 4. Inside the do simulation scope, you can also use the let function to define variables and bind them to simple values
See
Signature
declare const bindTo: {
<N extends string>(
name: N,
): <A, E, R>(self: Stream<A, E, R>) => Stream<{ [K in string]: A }, E, R>;
<A, E, R, N extends string>(self: Stream<A, E, R>, name: N): Stream<{ [K in string]: A }, E, R>;
};Example
import * as assert from "node:assert"
import { Chunk, Effect, pipe, Stream } from "effect"
const result = pipe(
Stream.Do,
Stream.bind("x", () => Stream.succeed(2)),
Stream.bind("y", () => Stream.succeed(3)),
Stream.let("sum", ({ x, y }) => x + y),
)
assert.deepStrictEqual(Effect.runSync(Stream.runCollect(result)), Chunk.of({ x: 2, y: 3, sum: 5 }))The "do simulation" in Effect allows you to write code in a more declarative style, similar to the "do notation" in other programming languages. It provides a way to define variables and perform operations on them using functions like bind and let.
Here's how the do simulation works:
1. Start the do simulation using the Do value 2. Within the do simulation scope, you can use the bind function to define variables and bind them to Stream values 3. You can accumulate multiple bind statements to define multiple variables within the scope 4. Inside the do simulation scope, you can also use the let function to define variables and bind them to simple values
See
Signature
declare const Do: Stream<{}>;Example
import * as assert from "node:assert"
import { Chunk, Effect, pipe, Stream } from "effect"
const result = pipe(
Stream.Do,
Stream.bind("x", () => Stream.succeed(2)),
Stream.bind("y", () => Stream.succeed(3)),
Stream.let("sum", ({ x, y }) => x + y),
)
assert.deepStrictEqual(Effect.runSync(Stream.runCollect(result)), Chunk.of({ x: 2, y: 3, sum: 5 }))Elements
Finds the first element emitted by this stream that satisfies the provided predicate.
Signature
declare const find: {
<A, B>(refinement: Refinement<NoInfer<A>, B>): <E, R>(self: Stream<A, E, R>) => Stream<B, E, R>;
<A>(predicate: Predicate<NoInfer<A>>): <E, R>(self: Stream<A, E, R>) => Stream<A, E, R>;
<A, E, R, B>(self: Stream<A, E, R>, refinement: Refinement<A, B>): Stream<B, E, R>;
<A, E, R>(self: Stream<A, E, R>, predicate: Predicate<A>): Stream<A, E, R>;
};findEffect
Finds the first element emitted by this stream that satisfies the provided effectful predicate.
Signature
declare const findEffect: {
<A, E2, R2>(
predicate: (a: NoInfer<A>) => Effect<boolean, E2, R2>,
): <E, R>(self: Stream<A, E, R>) => Stream<A, E2 | E, R2 | R>;
<A, E, R, E2, R2>(
self: Stream<A, E, R>,
predicate: (a: NoInfer<A>) => Effect<boolean, E2, R2>,
): Stream<A, E | E2, R | R2>;
};Encoding
decodeText
Decode Uint8Array chunks into a stream of strings using the specified encoding.
Signature
declare const decodeText: {
(
encoding?: string,
): <E, R>(self: Stream<Uint8Array<ArrayBufferLike>, E, R>) => Stream<string, E, R>;
<E, R>(self: Stream<Uint8Array<ArrayBufferLike>, E, R>, encoding?: string): Stream<string, E, R>;
};encodeText
Encode a stream of strings into a stream of Uint8Array chunks using the specified encoding.
Signature
declare const encodeText: <E, R>(self: Stream<string, E, R>) => Stream<Uint8Array, E, R>;Error Handling
Switches over to the stream produced by the provided function in case this one fails with a typed error.
Signature
declare const catchAll: {
<E, A2, E2, R2>(
f: (error: E) => Stream<A2, E2, R2>,
): <A, R>(self: Stream<A, E, R>) => Stream<A2 | A, E2, R2 | R>;
<A, E, R, A2, E2, R2>(
self: Stream<A, E, R>,
f: (error: E) => Stream<A2, E2, R2>,
): Stream<A | A2, E2, R | R2>;
};catchAllCause
Switches over to the stream produced by the provided function in case this one fails. Allows recovery from all causes of failure, including interruption if the stream is uninterruptible.
Signature
declare const catchAllCause: {
<E, A2, E2, R2>(
f: (cause: Cause<E>) => Stream<A2, E2, R2>,
): <A, R>(self: Stream<A, E, R>) => Stream<A2 | A, E2, R2 | R>;
<A, E, R, A2, E2, R2>(
self: Stream<A, E, R>,
f: (cause: Cause<E>) => Stream<A2, E2, R2>,
): Stream<A | A2, E2, R | R2>;
};Switches over to the stream produced by the provided function in case this one fails with some typed error.
Signature
declare const catchSome: {
<E, A2, E2, R2>(
pf: (error: E) => Option<Stream<A2, E2, R2>>,
): <A, R>(self: Stream<A, E, R>) => Stream<A2 | A, E | E2, R2 | R>;
<A, E, R, A2, E2, R2>(
self: Stream<A, E, R>,
pf: (error: E) => Option<Stream<A2, E2, R2>>,
): Stream<A | A2, E | E2, R | R2>;
};catchSomeCause
Switches over to the stream produced by the provided function in case this one fails with some errors. Allows recovery from all causes of failure, including interruption if the stream is uninterruptible.
Signature
declare const catchSomeCause: {
<E, A2, E2, R2>(
pf: (cause: Cause<E>) => Option<Stream<A2, E2, R2>>,
): <A, R>(self: Stream<A, E, R>) => Stream<A2 | A, E | E2, R2 | R>;
<A, E, R, A2, E2, R2>(
self: Stream<A, E, R>,
pf: (cause: Cause<E>) => Option<Stream<A2, E2, R2>>,
): Stream<A | A2, E | E2, R | R2>;
};Switches over to the stream produced by the provided function in case this one fails with an error matching the given _tag.
Signature
declare const catchTag: {
<
K extends string,
E extends {
_tag: string;
},
A1,
E1,
R1,
>(
k: K,
f: (
e: Extract<
E,
{
_tag: K;
}
>,
) => Stream<A1, E1, R1>,
): <A, R>(
self: Stream<A, E, R>,
) => Stream<
A1 | A,
| E1
| Exclude<
E,
{
_tag: K;
}
>,
R1 | R
>;
<
A,
E extends {
_tag: string;
},
R,
K extends string,
A1,
E1,
R1,
>(
self: Stream<A, E, R>,
k: K,
f: (
e: Extract<
E,
{
_tag: K;
}
>,
) => Stream<A1, E1, R1>,
): Stream<
A | A1,
| E1
| Exclude<
E,
{
_tag: K;
}
>,
R | R1
>;
};Switches over to the stream produced by one of the provided functions, in case this one fails with an error matching one of the given _tag's.
Signature
declare const catchTags: {
<
E extends {
_tag: string;
},
Cases extends {
[K in string]: (
error: Extract<
E,
{
_tag: K;
}
>,
) => Stream<any, any, any>;
},
>(
cases: Cases,
): <A, R>(
self: Stream<A, E, R>,
) => Stream<
| A
| {
[K in string | number | symbol]: Cases[K] extends (
...args: Array<any>
) => Variance<A, _E, _R>
? A
: never;
}[keyof Cases],
| Exclude<
E,
{
_tag: keyof Cases;
}
>
| {
[K in string | number | symbol]: Cases[K] extends (
...args: Array<any>
) => Variance<_A, E, _R>
? E
: never;
}[keyof Cases],
| R
| {
[K in string | number | symbol]: Cases[K] extends (
...args: Array<any>
) => Variance<_A, _E, R>
? R
: never;
}[keyof Cases]
>;
<
A,
E extends {
_tag: string;
},
R,
Cases extends {
[K in string]: (
error: Extract<
E,
{
_tag: K;
}
>,
) => Stream<any, any, any>;
},
>(
self: Stream<A, E, R>,
cases: Cases,
): Stream<
| A
| {
[K in string | number | symbol]: Cases[K] extends (
...args: Array<any>
) => Variance<_R, _E, A>
? A
: never;
}[keyof Cases],
| Exclude<
E,
{
_tag: keyof Cases;
}
>
| {
[K in string | number | symbol]: Cases[K] extends (
...args: Array<any>
) => Variance<_R, E, _A>
? E
: never;
}[keyof Cases],
| R
| {
[K in string | number | symbol]: Cases[K] extends (
...args: Array<any>
) => Variance<R, _E, _A>
? R
: never;
}[keyof Cases]
>;
};Translates any failure into a stream termination, making the stream infallible and all failures unchecked.
Signature
declare const orDie: <A, E, R>(self: Stream<A, E, R>) => Stream<A, never, R>;Keeps none of the errors, and terminates the stream with them, using the specified function to convert the E into a defect.
Signature
declare const orDieWith: {
<E>(f: (e: E) => unknown): <A, R>(self: Stream<A, E, R>) => Stream<A, never, R>;
<A, E, R>(self: Stream<A, E, R>, f: (e: E) => unknown): Stream<A, never, R>;
};Switches to the provided stream in case this one fails with a typed error.
See also Stream.catchAll.
Signature
declare const orElse: {
<A2, E2, R2>(
that: LazyArg<Stream<A2, E2, R2>>,
): <A, E, R>(self: Stream<A, E, R>) => Stream<A2 | A, E2, R2 | R>;
<A, E, R, A2, E2, R2>(
self: Stream<A, E, R>,
that: LazyArg<Stream<A2, E2, R2>>,
): Stream<A | A2, E2, R | R2>;
};orElseEither
Switches to the provided stream in case this one fails with a typed error.
See also Stream.catchAll.
Signature
declare const orElseEither: {
<A2, E2, R2>(
that: LazyArg<Stream<A2, E2, R2>>,
): <A, E, R>(self: Stream<A, E, R>) => Stream<Either<A2, A>, E2, R2 | R>;
<A, E, R, A2, E2, R2>(
self: Stream<A, E, R>,
that: LazyArg<Stream<A2, E2, R2>>,
): Stream<Either<A2, A>, E2, R | R2>;
};orElseFail
Fails with given error in case this one fails with a typed error.
See also Stream.catchAll.
Signature
declare const orElseFail: {
<E2>(error: LazyArg<E2>): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E2, R>;
<A, E, R, E2>(self: Stream<A, E, R>, error: LazyArg<E2>): Stream<A, E2, R>;
};orElseIfEmpty
Produces the specified element if this stream is empty.
Signature
declare const orElseIfEmpty: {
<A2>(element: LazyArg<A2>): <A, E, R>(self: Stream<A, E, R>) => Stream<A2 | A, E, R>;
<A, E, R, A2>(self: Stream<A, E, R>, element: LazyArg<A2>): Stream<A | A2, E, R>;
};orElseIfEmptyChunk
Produces the specified chunk if this stream is empty.
Signature
declare const orElseIfEmptyChunk: {
<A2>(chunk: LazyArg<Chunk<A2>>): <A, E, R>(self: Stream<A, E, R>) => Stream<A2 | A, E, R>;
<A, E, R, A2>(self: Stream<A, E, R>, chunk: LazyArg<Chunk<A2>>): Stream<A | A2, E, R>;
};orElseIfEmptyStream
Switches to the provided stream in case this one is empty.
Signature
declare const orElseIfEmptyStream: {
<A2, E2, R2>(
stream: LazyArg<Stream<A2, E2, R2>>,
): <A, E, R>(self: Stream<A, E, R>) => Stream<A2 | A, E2 | E, R2 | R>;
<A, E, R, A2, E2, R2>(
self: Stream<A, E, R>,
stream: LazyArg<Stream<A2, E2, R2>>,
): Stream<A | A2, E | E2, R | R2>;
};orElseSucceed
Succeeds with the specified value if this one fails with a typed error.
Signature
declare const orElseSucceed: {
<A2>(value: LazyArg<A2>): <A, E, R>(self: Stream<A, E, R>) => Stream<A2 | A, never, R>;
<A, E, R, A2>(self: Stream<A, E, R>, value: LazyArg<A2>): Stream<A | A2, never, R>;
};refineOrDie
Keeps some of the errors, and terminates the fiber with the rest
Signature
declare const refineOrDie: {
<E, E2>(pf: (error: E) => Option<E2>): <A, R>(self: Stream<A, E, R>) => Stream<A, E2, R>;
<A, E, R, E2>(self: Stream<A, E, R>, pf: (error: E) => Option<E2>): Stream<A, E2, R>;
};refineOrDieWith
Keeps some of the errors, and terminates the fiber with the rest, using the specified function to convert the E into a defect.
Signature
declare const refineOrDieWith: {
<E, E2>(
pf: (error: E) => Option<E2>,
f: (error: E) => unknown,
): <A, R>(self: Stream<A, E, R>) => Stream<A, E2, R>;
<A, E, R, E2>(
self: Stream<A, E, R>,
pf: (error: E) => Option<E2>,
f: (error: E) => unknown,
): Stream<A, E2, R>;
};Error Handling
withExecutionPlan
Apply an ExecutionPlan to the stream, which allows you to fallback to different resources in case of failure.
If you have a stream that could fail with partial results, you can use the preventFallbackOnPartialStream option to prevent contamination of the final stream with partial results.
Signature
declare const withExecutionPlan: {
<Input, R2, Provides, PolicyE>(
policy: ExecutionPlan<{
error: PolicyE;
input: Input;
provides: Provides;
requirements: R2;
}>,
options?: {
readonly preventFallbackOnPartialStream?: boolean;
},
): <A, E, R>(self: Stream<A, E, R>) => Stream<A, PolicyE | E, R2 | Exclude<R, Provides>>;
<A, E, R, R2, Input, Provides, PolicyE>(
self: Stream<A, E, R>,
policy: ExecutionPlan<{
error: PolicyE;
input: Input;
provides: Provides;
requirements: R2;
}>,
options?: {
readonly preventFallbackOnPartialStream?: boolean;
},
): Stream<A, E | PolicyE, R2 | Exclude<R, Provides>>;
};Filtering
Filters the elements emitted by this stream using the provided function.
Signature
declare const filter: {
<A, B>(refinement: Refinement<NoInfer<A>, B>): <E, R>(self: Stream<A, E, R>) => Stream<B, E, R>;
<A, B>(predicate: Predicate<B>): <E, R>(self: Stream<A, E, R>) => Stream<A, E, R>;
<A, E, R, B>(self: Stream<A, E, R>, refinement: Refinement<A, B>): Stream<B, E, R>;
<A, E, R>(self: Stream<A, E, R>, predicate: Predicate<A>): Stream<A, E, R>;
};Example
import { Effect, Stream } from "effect"
const stream = Stream.range(1, 11).pipe(Stream.filter((n) => n % 2 === 0))
Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// { _id: 'Chunk', values: [ 2, 4, 6, 8, 10 ] }filterEffect
Effectfully filters the elements emitted by this stream.
Signature
declare const filterEffect: {
<A, E2, R2>(
f: (a: NoInfer<A>) => Effect<boolean, E2, R2>,
): <E, R>(self: Stream<A, E, R>) => Stream<A, E2 | E, R2 | R>;
<A, E, R, E2, R2>(
self: Stream<A, E, R>,
f: (a: A) => Effect<boolean, E2, R2>,
): Stream<A, E | E2, R | R2>;
};Grouping
groupAdjacentBy
Creates a pipeline that groups on adjacent keys, calculated by the specified function.
Signature
declare const groupAdjacentBy: {
<A, K>(f: (a: A) => K): <E, R>(self: Stream<A, E, R>) => Stream<[K, NonEmptyChunk<A>], E, R>;
<A, E, R, K>(self: Stream<A, E, R>, f: (a: A) => K): Stream<[K, NonEmptyChunk<A>], E, R>;
};More powerful version of Stream.groupByKey.
Signature
declare const groupBy: {
<A, K, V, E2, R2>(
f: (a: A) => Effect<readonly [K, V], E2, R2>,
options?: {
readonly bufferSize?: number;
},
): <E, R>(self: Stream<A, E, R>) => GroupBy<K, V, E2 | E, R2 | R>;
<A, E, R, K, V, E2, R2>(
self: Stream<A, E, R>,
f: (a: A) => Effect<readonly [K, V], E2, R2>,
options?: {
readonly bufferSize?: number;
},
): GroupBy<K, V, E | E2, R | R2>;
};Example
import { Chunk, Effect, GroupBy, Stream } from "effect"
const groupByKeyResult = Stream.fromIterable([
"Mary",
"James",
"Robert",
"Patricia",
"John",
"Jennifer",
"Rebecca",
"Peter",
]).pipe(Stream.groupBy((name) => Effect.succeed([name.substring(0, 1), name])))
const stream = GroupBy.evaluate(groupByKeyResult, (key, stream) =>
Stream.fromEffect(
Stream.runCollect(stream).pipe(Effect.andThen((chunk) => [key, Chunk.size(chunk)] as const)),
),
)
Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// {
// _id: 'Chunk',
// values: [ [ 'M', 1 ], [ 'J', 3 ], [ 'R', 2 ], [ 'P', 2 ] ]
// }groupByKey
Partition a stream using a function and process each stream individually. This returns a data structure that can be used to further filter down which groups shall be processed.
After calling apply on the GroupBy object, the remaining groups will be processed in parallel and the resulting streams merged in a nondeterministic fashion.
Up to buffer elements may be buffered in any group stream before the producer is backpressured. Take care to consume from all streams in order to prevent deadlocks.
For example, to collect the first 2 words for every starting letter from a stream of words:
Signature
declare const groupByKey: {
<A, K>(
f: (a: A) => K,
options?: {
readonly bufferSize?: number;
},
): <E, R>(self: Stream<A, E, R>) => GroupBy<K, A, E, R>;
<A, E, R, K>(
self: Stream<A, E, R>,
f: (a: A) => K,
options?: {
readonly bufferSize?: number;
},
): GroupBy<K, A, E, R>;
};Example
import { pipe, GroupBy, Stream } from "effect"
pipe(
Stream.fromIterable(["hello", "world", "hi", "holla"]),
Stream.groupByKey((word) => word[0]),
GroupBy.evaluate((key, stream) =>
pipe(
stream,
Stream.take(2),
Stream.map((words) => [key, words] as const),
),
),
)Partitions the stream with specified chunkSize.
Signature
declare const grouped: {
(chunkSize: number): <A, E, R>(self: Stream<A, E, R>) => Stream<Chunk<A>, E, R>;
<A, E, R>(self: Stream<A, E, R>, chunkSize: number): Stream<Chunk<A>, E, R>;
};Example
import { Effect, Stream } from "effect"
const stream = Stream.range(0, 8).pipe(Stream.grouped(3))
Effect.runPromise(Stream.runCollect(stream)).then((chunks) => console.log("%o", chunks))
// {
// _id: 'Chunk',
// values: [
// { _id: 'Chunk', values: [ 0, 1, 2, [length]: 3 ] },
// { _id: 'Chunk', values: [ 3, 4, 5, [length]: 3 ] },
// { _id: 'Chunk', values: [ 6, 7, 8, [length]: 3 ] },
// [length]: 3
// ]
// }groupedWithin
Partitions the stream with the specified chunkSize or until the specified duration has passed, whichever is satisfied first.
Signature
declare const groupedWithin: {
(
chunkSize: number,
duration: DurationInput,
): <A, E, R>(self: Stream<A, E, R>) => Stream<Chunk<A>, E, R>;
<A, E, R>(
self: Stream<A, E, R>,
chunkSize: number,
duration: DurationInput,
): Stream<Chunk<A>, E, R>;
};Example
import { Chunk, Effect, Schedule, Stream } from "effect"
const stream = Stream.range(0, 9).pipe(
Stream.repeat(Schedule.spaced("1 second")),
Stream.groupedWithin(18, "1.5 seconds"),
Stream.take(3),
)
Effect.runPromise(Stream.runCollect(stream)).then((chunks) => console.log(Chunk.toArray(chunks)))
// [
// {
// _id: 'Chunk',
// values: [
// 0, 1, 2, 3, 4, 5, 6,
// 7, 8, 9, 0, 1, 2, 3,
// 4, 5, 6, 7
// ]
// },
// {
// _id: 'Chunk',
// values: [
// 8, 9, 0, 1, 2,
// 3, 4, 5, 6, 7,
// 8, 9
// ]
// },
// {
// _id: 'Chunk',
// values: [
// 0, 1, 2, 3, 4, 5, 6,
// 7, 8, 9, 0, 1, 2, 3,
// 4, 5, 6, 7
// ]
// }
// ]Mapping
Maps the success values of this stream to the specified constant value.
Signature
declare const as: {
<B>(value: B): <A, E, R>(self: Stream<A, E, R>) => Stream<B, E, R>;
<A, E, R, B>(self: Stream<A, E, R>, value: B): Stream<B, E, R>;
};Example
import { Effect, Stream } from "effect"
const stream = Stream.range(1, 5).pipe(Stream.as(null))
Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// { _id: 'Chunk', values: [ null, null, null, null, null ] }Transforms the elements of this stream using the supplied function.
Signature
declare const map: {
<A, B>(f: (a: A) => B): <E, R>(self: Stream<A, E, R>) => Stream<B, E, R>;
<A, E, R, B>(self: Stream<A, E, R>, f: (a: A) => B): Stream<B, E, R>;
};Example
import { Effect, Stream } from "effect"
const stream = Stream.make(1, 2, 3).pipe(Stream.map((n) => n + 1))
Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// { _id: 'Chunk', values: [ 2, 3, 4 ] }Statefully maps over the elements of this stream to produce new elements.
Signature
declare const mapAccum: {
<S, A, A2>(
s: S,
f: (s: S, a: A) => readonly [S, A2],
): <E, R>(self: Stream<A, E, R>) => Stream<A2, E, R>;
<A, E, R, S, A2>(
self: Stream<A, E, R>,
s: S,
f: (s: S, a: A) => readonly [S, A2],
): Stream<A2, E, R>;
};Example
import { Effect, Stream } from "effect"
const runningTotal = (stream: Stream.Stream<number>): Stream.Stream<number> =>
stream.pipe(Stream.mapAccum(0, (s, a) => [s + a, s + a]))
// input: 0, 1, 2, 3, 4, 5, 6
Effect.runPromise(Stream.runCollect(runningTotal(Stream.range(0, 6)))).then(console.log)
// { _id: "Chunk", values: [ 0, 1, 3, 6, 10, 15, 21 ] }mapAccumEffect
Statefully and effectfully maps over the elements of this stream to produce new elements.
Signature
declare const mapAccumEffect: {
<S, A, A2, E2, R2>(
s: S,
f: (s: S, a: A) => Effect<readonly [S, A2], E2, R2>,
): <E, R>(self: Stream<A, E, R>) => Stream<A2, E2 | E, R2 | R>;
<A, E, R, S, A2, E2, R2>(
self: Stream<A, E, R>,
s: S,
f: (s: S, a: A) => Effect<readonly [S, A2], E2, R2>,
): Stream<A2, E | E2, R | R2>;
};Transforms the chunks emitted by this stream.
Signature
declare const mapChunks: {
<A, B>(f: (chunk: Chunk<A>) => Chunk<B>): <E, R>(self: Stream<A, E, R>) => Stream<B, E, R>;
<A, E, R, B>(self: Stream<A, E, R>, f: (chunk: Chunk<A>) => Chunk<B>): Stream<B, E, R>;
};mapChunksEffect
Effectfully transforms the chunks emitted by this stream.
Signature
declare const mapChunksEffect: {
<A, B, E2, R2>(
f: (chunk: Chunk<A>) => Effect<Chunk<B>, E2, R2>,
): <E, R>(self: Stream<A, E, R>) => Stream<B, E2 | E, R2 | R>;
<A, E, R, B, E2, R2>(
self: Stream<A, E, R>,
f: (chunk: Chunk<A>) => Effect<Chunk<B>, E2, R2>,
): Stream<B, E | E2, R | R2>;
};Maps each element to an iterable, and flattens the iterables into the output of this stream.
Signature
declare const mapConcat: {
<A, A2>(f: (a: A) => Iterable<A2>): <E, R>(self: Stream<A, E, R>) => Stream<A2, E, R>;
<A, E, R, A2>(self: Stream<A, E, R>, f: (a: A) => Iterable<A2>): Stream<A2, E, R>;
};Example
import { Effect, Stream } from "effect"
const numbers = Stream.make("1-2-3", "4-5", "6").pipe(
Stream.mapConcat((s) => s.split("-")),
Stream.map((s) => parseInt(s)),
)
Effect.runPromise(Stream.runCollect(numbers)).then(console.log)
// { _id: 'Chunk', values: [ 1, 2, 3, 4, 5, 6 ] }mapConcatChunk
Maps each element to a chunk, and flattens the chunks into the output of this stream.
Signature
declare const mapConcatChunk: {
<A, A2>(f: (a: A) => Chunk<A2>): <E, R>(self: Stream<A, E, R>) => Stream<A2, E, R>;
<A, E, R, A2>(self: Stream<A, E, R>, f: (a: A) => Chunk<A2>): Stream<A2, E, R>;
};mapConcatChunkEffect
Effectfully maps each element to a chunk, and flattens the chunks into the output of this stream.
Signature
declare const mapConcatChunkEffect: {
<A, A2, E2, R2>(
f: (a: A) => Effect<Chunk<A2>, E2, R2>,
): <E, R>(self: Stream<A, E, R>) => Stream<A2, E2 | E, R2 | R>;
<A, E, R, A2, E2, R2>(
self: Stream<A, E, R>,
f: (a: A) => Effect<Chunk<A2>, E2, R2>,
): Stream<A2, E | E2, R | R2>;
};mapConcatEffect
Effectfully maps each element to an iterable, and flattens the iterables into the output of this stream.
Signature
declare const mapConcatEffect: {
<A, A2, E2, R2>(
f: (a: A) => Effect<Iterable<A2, any, any>, E2, R2>,
): <E, R>(self: Stream<A, E, R>) => Stream<A2, E2 | E, R2 | R>;
<A, E, R, A2, E2, R2>(
self: Stream<A, E, R>,
f: (a: A) => Effect<Iterable<A2, any, any>, E2, R2>,
): Stream<A2, E | E2, R | R2>;
};Maps over elements of the stream with the specified effectful function.
Signature
declare const mapEffect: {
<A, A2, E2, R2>(
f: (a: A) => Effect<A2, E2, R2>,
options?: {
readonly concurrency?: number | "unbounded";
readonly unordered?: boolean;
},
): <E, R>(self: Stream<A, E, R>) => Stream<A2, E2 | E, R2 | R>;
<A, A2, E2, R2, K>(
f: (a: A) => Effect<A2, E2, R2>,
options: {
readonly bufferSize?: number;
readonly key: (a: A) => K;
},
): <E, R>(self: Stream<A, E, R>) => Stream<A2, E2 | E, R2 | R>;
<A, E, R, A2, E2, R2>(
self: Stream<A, E, R>,
f: (a: A) => Effect<A2, E2, R2>,
options?: {
readonly concurrency?: number | "unbounded";
readonly unordered?: boolean;
},
): Stream<A2, E | E2, R | R2>;
<A, E, R, A2, E2, R2, K>(
self: Stream<A, E, R>,
f: (a: A) => Effect<A2, E2, R2>,
options: {
readonly bufferSize?: number;
readonly key: (a: A) => K;
},
): Stream<A2, E | E2, R | R2>;
};Example
import { Effect, Random, Stream } from "effect"
const stream = Stream.make(10, 20, 30).pipe(Stream.mapEffect((n) => Random.nextIntBetween(0, n)))
Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// Example Output: { _id: 'Chunk', values: [ 7, 19, 8 ] }Transforms the errors emitted by this stream using f.
Signature
declare const mapError: {
<E, E2>(f: (error: E) => E2): <A, R>(self: Stream<A, E, R>) => Stream<A, E2, R>;
<A, E, R, E2>(self: Stream<A, E, R>, f: (error: E) => E2): Stream<A, E2, R>;
};mapErrorCause
Transforms the full causes of failures emitted by this stream.
Signature
declare const mapErrorCause: {
<E, E2>(f: (cause: Cause<E>) => Cause<E2>): <A, R>(self: Stream<A, E, R>) => Stream<A, E2, R>;
<A, E, R, E2>(self: Stream<A, E, R>, f: (cause: Cause<E>) => Cause<E2>): Stream<A, E2, R>;
};Models
EventListener interface
Signature
interface EventListener<A> {
addEventListener(
event: string,
f: (event: A) => void,
options?:
| boolean
| {
readonly capture?: boolean;
readonly once?: boolean;
readonly passive?: boolean;
readonly signal?: AbortSignal;
},
): void;
removeEventListener(
event: string,
f: (event: A) => void,
options?:
| boolean
| {
readonly capture?: boolean;
},
): void;
}A Stream<A, E, R> is a description of a program that, when evaluated, may emit zero or more values of type A, may fail with errors of type E, and uses an context of type R. One way to think of Stream is as a Effect program that could emit multiple values.
Stream is a purely functional *pull* based stream. Pull based streams offer inherent laziness and backpressure, relieving users of the need to manage buffers between operators. As an optimization, Stream does not emit single values, but rather an array of values. This allows the cost of effect evaluation to be amortized.
Stream forms a monad on its A type parameter, and has error management facilities for its E type parameter, modeled similarly to Effect (with some adjustments for the multiple-valued nature of Stream). These aspects allow for rich and expressive composition of streams.
Signature
interface Stream<out A, out E = never, out R = never> extends Variance<A, E, R>, Pipeable {
[ignoreSymbol]?: StreamUnifyIgnore;
[typeSymbol]?: unknown;
[unifySymbol]?: StreamUnify<Stream<A, E, R>>;
}StreamUnify interface
Signature
interface StreamUnify<
A extends {
[typeSymbol]?: any;
},
> extends EffectUnify<A> {
Stream?: () => A[typeof typeSymbol] extends Stream<A0, E0, R0> | _ ? Stream<A0, E0, R0> : never;
}StreamUnifyIgnore interface
Signature
interface StreamUnifyIgnore extends EffectUnifyIgnore {
Effect?: true;
}Other
Signature
declare const async: <A, E = never, R = never>(
register: (emit: Emit.Emit<R, E, A, void>) => Effect.Effect<void, never, R> | void,
bufferSize?:
| number
| "unbounded"
| {
readonly bufferSize?: number;
readonly strategy?: "dropping" | "sliding" | "suspend";
},
) => Stream<A, E, R>;fromEventListener
Creates a Stream using addEventListener.
Signature
declare const fromEventListener: <A = unknown>(
target: EventListener<A>,
type: string,
options?:
| boolean
| {
readonly bufferSize?: number | "unbounded";
readonly capture?: boolean;
readonly once?: boolean;
readonly passive?: boolean;
},
) => Stream<A>;Signature
declare const let: {
<N extends string, A extends object, B>(
name: Exclude<N, keyof A>,
f: (a: NoInfer<A>) => B,
): <E, R>(
self: Stream<A, E, R>,
) => Stream<{ [K in string | number | symbol]: K extends keyof A ? A[K] : B }, E, R>;
<A extends object, E, R, N extends string, B>(
self: Stream<A, E, R>,
name: Exclude<N, keyof A>,
f: (a: NoInfer<A>) => B,
): Stream<{ [K in string | number | symbol]: K extends keyof A ? A[K] : B }, E, R>;
};Signature
declare const void: Stream<void>Racing
Returns a stream that mirrors the first upstream to emit an item. As soon as one of the upstream emits a first value, the other is interrupted. The resulting stream will forward all items from the "winning" source stream. Any upstream failures will cause the returned stream to fail.
Signature
declare const race: {
<AR, ER, RR>(
right: Stream<AR, ER, RR>,
): <AL, EL, RL>(left: Stream<AL, EL, RL>) => Stream<AR | AL, ER | EL, RR | RL>;
<AL, EL, RL, AR, ER, RR>(
left: Stream<AL, EL, RL>,
right: Stream<AR, ER, RR>,
): Stream<AL | AR, EL | ER, RL | RR>;
};Example
import { Stream, Schedule, Console, Effect } from "effect"
const stream = Stream.fromSchedule(Schedule.spaced("2 millis")).pipe(
Stream.race(Stream.fromSchedule(Schedule.spaced("1 millis"))),
Stream.take(6),
Stream.tap(Console.log),
)
Effect.runPromise(Stream.runDrain(stream))
// Output each millisecond from the first stream, the rest streams are interrupted
// 0
// 1
// 2
// 3
// 4
// 5Returns a stream that mirrors the first upstream to emit an item. As soon as one of the upstream emits a first value, all the others are interrupted. The resulting stream will forward all items from the "winning" source stream. Any upstream failures will cause the returned stream to fail.
Signature
declare const raceAll: <S extends ReadonlyArray<Stream<any, any, any>>>(
...streams: S
) => Stream<Stream.Success<S[number]>, Stream.Error<S[number]>, Stream.Context<S[number]>>;Example
import { Stream, Schedule, Console, Effect } from "effect"
const stream = Stream.raceAll(
Stream.fromSchedule(Schedule.spaced("1 millis")),
Stream.fromSchedule(Schedule.spaced("2 millis")),
Stream.fromSchedule(Schedule.spaced("4 millis")),
).pipe(Stream.take(6), Stream.tap(Console.log))
Effect.runPromise(Stream.runDrain(stream))
// Output each millisecond from the first stream, the rest streams are interrupted
// 0
// 1
// 2
// 3
// 4
// 5Sequencing
branchAfter
Returns a Stream that first collects n elements from the input Stream, and then creates a new Stream using the specified function, and sends all the following elements through that.
Signature
declare const branchAfter: {
<A, A2, E2, R2>(
n: number,
f: (input: Chunk<A>) => Stream<A2, E2, R2>,
): <E, R>(self: Stream<A, E, R>) => Stream<A2, E2 | E, R2 | R>;
<A, E, R, A2, E2, R2>(
self: Stream<A, E, R>,
n: number,
f: (input: Chunk<A>) => Stream<A2, E2, R2>,
): Stream<A2, E | E2, R | R2>;
};Returns a stream made of the concatenation in strict order of all the streams produced by passing each element of this stream to f0
Signature
declare const flatMap: {
<A, A2, E2, R2>(
f: (a: A) => Stream<A2, E2, R2>,
options?: {
readonly bufferSize?: number;
readonly concurrency?: number | "unbounded";
readonly switch?: boolean;
},
): <E, R>(self: Stream<A, E, R>) => Stream<A2, E2 | E, R2 | R>;
<A, E, R, A2, E2, R2>(
self: Stream<A, E, R>,
f: (a: A) => Stream<A2, E2, R2>,
options?: {
readonly bufferSize?: number;
readonly concurrency?: number | "unbounded";
readonly switch?: boolean;
},
): Stream<A2, E | E2, R | R2>;
};Flattens this stream-of-streams into a stream made of the concatenation in strict order of all the streams.
Signature
declare const flatten: {
(options?: {
readonly bufferSize?: number;
readonly concurrency?: number | "unbounded";
}): <A, E2, R2, E, R>(self: Stream<Stream<A, E2, R2>, E, R>) => Stream<A, E2 | E, R2 | R>;
<A, E2, R2, E, R>(
self: Stream<Stream<A, E2, R2>, E, R>,
options?: {
readonly bufferSize?: number;
readonly concurrency?: number | "unbounded";
},
): Stream<A, E2 | E, R2 | R>;
};flattenChunks
Submerges the chunks carried by this stream into the stream's structure, while still preserving them.
Signature
declare const flattenChunks: <A, E, R>(self: Stream<Chunk.Chunk<A>, E, R>) => Stream<A, E, R>;flattenEffect
Flattens Effect values into the stream's structure, preserving all information about the effect.
Signature
declare const flattenEffect: {
(options?: {
readonly concurrency?: number | "unbounded";
readonly unordered?: boolean;
}): <A, E2, R2, E, R>(self: Stream<Effect<A, E2, R2>, E, R>) => Stream<A, E2 | E, R2 | R>;
<A, E2, R2, E, R>(
self: Stream<Effect<A, E2, R2>, E, R>,
options?: {
readonly concurrency?: number | "unbounded";
readonly unordered?: boolean;
},
): Stream<A, E2 | E, R2 | R>;
};flattenExitOption
Unwraps Exit values that also signify end-of-stream by failing with None.
Signature
declare const flattenExitOption: <A, E2, E, R>(
self: Stream<Exit.Exit<A, Option.Option<E2>>, E, R>,
) => Stream<A, E | E2, R>;flattenIterables
Submerges the iterables carried by this stream into the stream's structure, while still preserving them.
Signature
declare const flattenIterables: <A, E, R>(self: Stream<Iterable<A>, E, R>) => Stream<A, E, R>;flattenTake
Unwraps Exit values and flatten chunks that also signify end-of-stream by failing with None.
Signature
declare const flattenTake: <A, E2, E, R>(
self: Stream<Take.Take<A, E2>, E, R>,
) => Stream<A, E | E2, R>;Adds an effect to be executed at the end of the stream.
Signature
declare const onEnd: {
<_, E2, R2>(
effect: Effect<_, E2, R2>,
): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E2 | E, R2 | R>;
<A, E, R, _, E2, R2>(self: Stream<A, E, R>, effect: Effect<_, E2, R2>): Stream<A, E | E2, R | R2>;
};Example
import { Console, Effect, Stream } from "effect"
const stream = Stream.make(1, 2, 3).pipe(
Stream.map((n) => n * 2),
Stream.tap((n) => Console.log(`after mapping: ${n}`)),
Stream.onEnd(Console.log("Stream ended")),
)
Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// after mapping: 2
// after mapping: 4
// after mapping: 6
// Stream ended
// { _id: 'Chunk', values: [ 2, 4, 6 ] }Adds an effect to be executed at the start of the stream.
Signature
declare const onStart: {
<_, E2, R2>(
effect: Effect<_, E2, R2>,
): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E2 | E, R2 | R>;
<A, E, R, _, E2, R2>(self: Stream<A, E, R>, effect: Effect<_, E2, R2>): Stream<A, E | E2, R | R2>;
};Example
import { Console, Effect, Stream } from "effect"
const stream = Stream.make(1, 2, 3).pipe(
Stream.onStart(Console.log("Stream started")),
Stream.map((n) => n * 2),
Stream.tap((n) => Console.log(`after mapping: ${n}`)),
)
Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// Stream started
// after mapping: 2
// after mapping: 4
// after mapping: 6
// { _id: 'Chunk', values: [ 2, 4, 6 ] }Adds an effect to consumption of every element of the stream.
Signature
declare const tap: {
<A, X, E2, R2>(
f: (a: NoInfer<A>) => Effect<X, E2, R2>,
): <E, R>(self: Stream<A, E, R>) => Stream<A, E2 | E, R2 | R>;
<A, E, R, X, E2, R2>(
self: Stream<A, E, R>,
f: (a: NoInfer<A>) => Effect<X, E2, R2>,
): Stream<A, E | E2, R | R2>;
};Example
import { Console, Effect, Stream } from "effect"
const stream = Stream.make(1, 2, 3).pipe(
Stream.tap((n) => Console.log(`before mapping: ${n}`)),
Stream.map((n) => n * 2),
Stream.tap((n) => Console.log(`after mapping: ${n}`)),
)
Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// before mapping: 1
// after mapping: 2
// before mapping: 2
// after mapping: 4
// before mapping: 3
// after mapping: 6
// { _id: 'Chunk', values: [ 2, 4, 6 ] }Returns a stream that effectfully "peeks" at the failure or success of the stream.
Signature
declare const tapBoth: {
<E, X1, E2, R2, A, X2, E3, R3>(options: {
readonly onFailure: (e: NoInfer<E>) => Effect.Effect<X1, E2, R2>;
readonly onSuccess: (a: NoInfer<A>) => Effect.Effect<X2, E3, R3>;
}): <R>(self: Stream<A, E, R>) => Stream<A, E | E2 | E3, R2 | R3 | R>;
<A, E, R, X1, E2, R2, X2, E3, R3>(
self: Stream<A, E, R>,
options: {
readonly onFailure: (e: NoInfer<E>) => Effect.Effect<X1, E2, R2>;
readonly onSuccess: (a: NoInfer<A>) => Effect.Effect<X2, E3, R3>;
},
): Stream<A, E | E2 | E3, R | R2 | R3>;
};Returns a stream that effectfully "peeks" at the failure of the stream.
Signature
declare const tapError: {
<E, X, E2, R2>(
f: (error: NoInfer<E>) => Effect<X, E2, R2>,
): <A, R>(self: Stream<A, E, R>) => Stream<A, E | E2, R2 | R>;
<A, E, R, X, E2, R2>(
self: Stream<A, E, R>,
f: (error: E) => Effect<X, E2, R2>,
): Stream<A, E | E2, R | R2>;
};Sends all elements emitted by this stream to the specified sink in addition to emitting them.
Signature
declare const tapSink: {
<A, E2, R2>(
sink: Sink<unknown, A, unknown, E2, R2>,
): <E, R>(self: Stream<A, E, R>) => Stream<A, E2 | E, R2 | R>;
<A, E, R, E2, R2>(
self: Stream<A, E, R>,
sink: Sink<unknown, A, unknown, E2, R2>,
): Stream<A, E | E2, R | R2>;
};Symbols
StreamTypeId
Signature
declare const StreamTypeId: unique symbol;StreamTypeId type
Signature
type StreamTypeId = typeof StreamTypeId;Tracing
Wraps the stream with a new span for tracing.
Signature
declare const withSpan: {
(
name: string,
options?: SpanOptions,
): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E, Exclude<R, ParentSpan>>;
<A, E, R>(
self: Stream<A, E, R>,
name: string,
options?: SpanOptions,
): Stream<A, E, Exclude<R, ParentSpan>>;
};Type Lambdas
StreamTypeLambda interface
Signature
interface StreamTypeLambda extends TypeLambda {
readonly type: Stream<unknown, unknown, unknown>;
}Utils
accumulate
Collects each underlying Chunk of the stream into a new chunk, and emits it on each pull.
Signature
declare const accumulate: <A, E, R>(self: Stream<A, E, R>) => Stream<Chunk.Chunk<A>, E, R>;accumulateChunks
Re-chunks the elements of the stream by accumulating each underlying chunk.
Signature
declare const accumulateChunks: <A, E, R>(self: Stream<A, E, R>) => Stream<A, E, R>;Aggregates elements of this stream using the provided sink for as long as the downstream operators on the stream are busy.
This operator divides the stream into two asynchronous "islands". Operators upstream of this operator run on one fiber, while downstream operators run on another. Whenever the downstream fiber is busy processing elements, the upstream fiber will feed elements into the sink until it signals completion.
Any sink can be used here, but see Sink.foldWeightedEffect and Sink.foldUntilEffect for sinks that cover the common usecases.
Signature
declare const aggregate: {
<B, A, A2, E2, R2>(
sink: Sink<B, A | A2, A2, E2, R2>,
): <E, R>(self: Stream<A, E, R>) => Stream<B, E2 | E, R2 | R>;
<A, E, R, B, A2, E2, R2>(
self: Stream<A, E, R>,
sink: Sink<B, A | A2, A2, E2, R2>,
): Stream<B, E | E2, R | R2>;
};aggregateWithin
Like aggregateWithinEither, but only returns the Right results.
Signature
declare const aggregateWithin: {
<B, A, A2, E2, R2, C, R3>(
sink: Sink<B, A | A2, A2, E2, R2>,
schedule: Schedule<C, Option<B>, R3>,
): <E, R>(self: Stream<A, E, R>) => Stream<B, E2 | E, R2 | R3 | R>;
<A, E, R, B, A2, E2, R2, C, R3>(
self: Stream<A, E, R>,
sink: Sink<B, A | A2, A2, E2, R2>,
schedule: Schedule<C, Option<B>, R3>,
): Stream<B, E | E2, R | R2 | R3>;
};aggregateWithinEither
Aggregates elements using the provided sink until it completes, or until the delay signalled by the schedule has passed.
This operator divides the stream into two asynchronous islands. Operators upstream of this operator run on one fiber, while downstream operators run on another. Elements will be aggregated by the sink until the downstream fiber pulls the aggregated value, or until the schedule's delay has passed.
Aggregated elements will be fed into the schedule to determine the delays between pulls.
Signature
declare const aggregateWithinEither: {
<B, A, A2, E2, R2, C, R3>(
sink: Sink<B, A | A2, A2, E2, R2>,
schedule: Schedule<C, Option<B>, R3>,
): <E, R>(self: Stream<A, E, R>) => Stream<Either<B, C>, E2 | E, R2 | R3 | R>;
<A, E, R, B, A2, E2, R2, C, R3>(
self: Stream<A, E, R>,
sink: Sink<B, A | A2, A2, E2, R2>,
schedule: Schedule<C, Option<B>, R3>,
): Stream<Either<B, C>, E | E2, R | R2 | R3>;
};Fan out the stream, producing a list of streams that have the same elements as this stream. The driver stream will only ever advance the maximumLag chunks before the slowest downstream stream.
Signature
declare const broadcast: {
<N extends number>(
n: N,
maximumLag:
| number
| {
readonly capacity: "unbounded";
readonly replay?: number;
}
| {
readonly capacity: number;
readonly replay?: number;
readonly strategy?: "sliding" | "dropping" | "suspend";
},
): <A, E, R>(self: Stream<A, E, R>) => Effect<TupleOf<N, Stream<A, E, never>>, never, Scope | R>;
<A, E, R, N extends number>(
self: Stream<A, E, R>,
n: N,
maximumLag:
| number
| {
readonly capacity: "unbounded";
readonly replay?: number;
}
| {
readonly capacity: number;
readonly replay?: number;
readonly strategy?: "sliding" | "dropping" | "suspend";
},
): Effect<TupleOf<N, Stream<A, E, never>>, never, Scope | R>;
};Example
import { Console, Effect, Fiber, Schedule, Stream } from "effect"
const numbers = Effect.scoped(
Stream.range(1, 20).pipe(
Stream.tap((n) => Console.log(`Emit ${n} element before broadcasting`)),
Stream.broadcast(2, 5),
Stream.flatMap(([first, second]) =>
Effect.gen(function* () {
const fiber1 = yield* Stream.runFold(first, 0, (acc, e) => Math.max(acc, e)).pipe(
Effect.andThen((max) => Console.log(`Maximum: ${max}`)),
Effect.fork,
)
const fiber2 = yield* second.pipe(
Stream.schedule(Schedule.spaced("1 second")),
Stream.runForEach((n) => Console.log(`Logging to the Console: ${n}`)),
Effect.fork,
)
yield* Fiber.join(fiber1).pipe(Effect.zip(Fiber.join(fiber2), { concurrent: true }))
}),
),
Stream.runCollect,
),
)
Effect.runPromise(numbers).then(console.log)
// Emit 1 element before broadcasting
// Emit 2 element before broadcasting
// Emit 3 element before broadcasting
// Emit 4 element before broadcasting
// Emit 5 element before broadcasting
// Emit 6 element before broadcasting
// Emit 7 element before broadcasting
// Emit 8 element before broadcasting
// Emit 9 element before broadcasting
// Emit 10 element before broadcasting
// Emit 11 element before broadcasting
// Logging to the Console: 1
// Logging to the Console: 2
// Logging to the Console: 3
// Logging to the Console: 4
// Logging to the Console: 5
// Emit 12 element before broadcasting
// Emit 13 element before broadcasting
// Emit 14 element before broadcasting
// Emit 15 element before broadcasting
// Emit 16 element before broadcasting
// Logging to the Console: 6
// Logging to the Console: 7
// Logging to the Console: 8
// Logging to the Console: 9
// Logging to the Console: 10
// Emit 17 element before broadcasting
// Emit 18 element before broadcasting
// Emit 19 element before broadcasting
// Emit 20 element before broadcasting
// Logging to the Console: 11
// Logging to the Console: 12
// Logging to the Console: 13
// Logging to the Console: 14
// Logging to the Console: 15
// Maximum: 20
// Logging to the Console: 16
// Logging to the Console: 17
// Logging to the Console: 18
// Logging to the Console: 19
// Logging to the Console: 20
// { _id: 'Chunk', values: [ undefined ] }broadcastDynamic
Fan out the stream, producing a dynamic number of streams that have the same elements as this stream. The driver stream will only ever advance the maximumLag chunks before the slowest downstream stream.
Signature
declare const broadcastDynamic: {
(
maximumLag:
| number
| {
readonly capacity: "unbounded";
readonly replay?: number;
}
| {
readonly capacity: number;
readonly replay?: number;
readonly strategy?: "sliding" | "dropping" | "suspend";
},
): <A, E, R>(self: Stream<A, E, R>) => Effect<Stream<A, E, never>, never, Scope | R>;
<A, E, R>(
self: Stream<A, E, R>,
maximumLag:
| number
| {
readonly capacity: "unbounded";
readonly replay?: number;
}
| {
readonly capacity: number;
readonly replay?: number;
readonly strategy?: "sliding" | "dropping" | "suspend";
},
): Effect<Stream<A, E, never>, never, Scope | R>;
};broadcastedQueues
Converts the stream to a scoped list of queues. Every value will be replicated to every queue with the slowest queue being allowed to buffer maximumLag chunks before the driver is back pressured.
Queues can unsubscribe from upstream by shutting down.
Signature
declare const broadcastedQueues: {
<N extends number>(
n: N,
maximumLag:
| number
| {
readonly capacity: "unbounded";
readonly replay?: number;
}
| {
readonly capacity: number;
readonly replay?: number;
readonly strategy?: "sliding" | "dropping" | "suspend";
},
): <A, E, R>(self: Stream<A, E, R>) => Effect<TupleOf<N, Dequeue<Take<A, E>>>, never, Scope | R>;
<A, E, R, N extends number>(
self: Stream<A, E, R>,
n: N,
maximumLag:
| number
| {
readonly capacity: "unbounded";
readonly replay?: number;
}
| {
readonly capacity: number;
readonly replay?: number;
readonly strategy?: "sliding" | "dropping" | "suspend";
},
): Effect<TupleOf<N, Dequeue<Take<A, E>>>, never, Scope | R>;
};broadcastedQueuesDynamic
Converts the stream to a scoped dynamic amount of queues. Every chunk will be replicated to every queue with the slowest queue being allowed to buffer maximumLag chunks before the driver is back pressured.
Queues can unsubscribe from upstream by shutting down.
Signature
declare const broadcastedQueuesDynamic: {
(
maximumLag:
| number
| {
readonly capacity: "unbounded";
readonly replay?: number;
}
| {
readonly capacity: number;
readonly replay?: number;
readonly strategy?: "sliding" | "dropping" | "suspend";
},
): <A, E, R>(
self: Stream<A, E, R>,
) => Effect<Effect<Dequeue<Take<A, E>>, never, Scope>, never, Scope | R>;
<A, E, R>(
self: Stream<A, E, R>,
maximumLag:
| number
| {
readonly capacity: "unbounded";
readonly replay?: number;
}
| {
readonly capacity: number;
readonly replay?: number;
readonly strategy?: "sliding" | "dropping" | "suspend";
},
): Effect<Effect<Dequeue<Take<A, E>>, never, Scope>, never, Scope | R>;
};Allows a faster producer to progress independently of a slower consumer by buffering up to capacity elements in a queue.
Note: This combinator destroys the chunking structure. It's recommended to use rechunk afterwards. Additionally, prefer capacities that are powers of 2 for better performance.
Signature
declare const buffer: {
(
options:
| {
readonly capacity: "unbounded";
}
| {
readonly capacity: number;
readonly strategy?: "dropping" | "sliding" | "suspend";
},
): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E, R>;
<A, E, R>(
self: Stream<A, E, R>,
options:
| {
readonly capacity: "unbounded";
}
| {
readonly capacity: number;
readonly strategy?: "dropping" | "sliding" | "suspend";
},
): Stream<A, E, R>;
};Example
import { Console, Effect, Schedule, Stream } from "effect"
const stream = Stream.range(1, 10).pipe(
Stream.tap((n) => Console.log(`before buffering: ${n}`)),
Stream.buffer({ capacity: 4 }),
Stream.tap((n) => Console.log(`after buffering: ${n}`)),
Stream.schedule(Schedule.spaced("5 seconds")),
)
Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// before buffering: 1
// before buffering: 2
// before buffering: 3
// before buffering: 4
// before buffering: 5
// before buffering: 6
// after buffering: 1
// after buffering: 2
// before buffering: 7
// after buffering: 3
// before buffering: 8
// after buffering: 4
// before buffering: 9
// after buffering: 5
// before buffering: 10
// ...bufferChunks
Allows a faster producer to progress independently of a slower consumer by buffering up to capacity chunks in a queue.
Signature
declare const bufferChunks: {
(options: {
readonly capacity: number;
readonly strategy?: "dropping" | "sliding" | "suspend";
}): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E, R>;
<A, E, R>(
self: Stream<A, E, R>,
options: {
readonly capacity: number;
readonly strategy?: "dropping" | "sliding" | "suspend";
},
): Stream<A, E, R>;
};Returns a new stream that only emits elements that are not equal to the previous element emitted, using natural equality to determine whether two elements are equal.
Signature
declare const changes: <A, E, R>(self: Stream<A, E, R>) => Stream<A, E, R>;Example
import { Effect, Stream } from "effect"
const stream = Stream.make(1, 1, 1, 2, 2, 3, 4).pipe(Stream.changes)
Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// { _id: 'Chunk', values: [ 1, 2, 3, 4 ] }changesWith
Returns a new stream that only emits elements that are not equal to the previous element emitted, using the specified function to determine whether two elements are equal.
Signature
declare const changesWith: {
<A>(f: (x: A, y: A) => boolean): <E, R>(self: Stream<A, E, R>) => Stream<A, E, R>;
<A, E, R>(self: Stream<A, E, R>, f: (x: A, y: A) => boolean): Stream<A, E, R>;
};changesWithEffect
Returns a new stream that only emits elements that are not equal to the previous element emitted, using the specified effectual function to determine whether two elements are equal.
Signature
declare const changesWithEffect: {
<A, E2, R2>(
f: (x: A, y: A) => Effect<boolean, E2, R2>,
): <E, R>(self: Stream<A, E, R>) => Stream<A, E2 | E, R2 | R>;
<A, E, R, E2, R2>(
self: Stream<A, E, R>,
f: (x: A, y: A) => Effect<boolean, E2, R2>,
): Stream<A, E | E2, R | R2>;
};Exposes the underlying chunks of the stream as a stream of chunks of elements.
Signature
declare const chunks: <A, E, R>(self: Stream<A, E, R>) => Stream<Chunk.Chunk<A>, E, R>;chunksWith
Performs the specified stream transformation with the chunk structure of the stream exposed.
Signature
declare const chunksWith: {
<A, E, R, A2, E2, R2>(
f: (stream: Stream<Chunk<A>, E, R>) => Stream<Chunk<A2>, E2, R2>,
): (self: Stream<A, E, R>) => Stream<A2, E | E2, R | R2>;
<A, E, R, A2, E2, R2>(
self: Stream<A, E, R>,
f: (stream: Stream<Chunk<A>, E, R>) => Stream<Chunk<A2>, E2, R2>,
): Stream<A2, E | E2, R | R2>;
};Combines the elements from this stream and the specified stream by repeatedly applying the function f to extract an element using both sides and conceptually "offer" it to the destination stream. f can maintain some internal state to control the combining process, with the initial state being specified by s.
Where possible, prefer Stream.combineChunks for a more efficient implementation.
Signature
declare const combine: {
<A2, E2, R2, S, R3, E, A, R4, R5, A3>(
that: Stream<A2, E2, R2>,
s: S,
f: (
s: S,
pullLeft: Effect<A, Option<E>, R3>,
pullRight: Effect<A2, Option<E2>, R4>,
) => Effect<Exit<readonly [A3, S], Option<E2 | E>>, never, R5>,
): <R>(self: Stream<A, E, R>) => Stream<A3, E2 | E, R2 | R3 | R4 | R5 | R>;
<R, A2, E2, R2, S, R3, E, A, R4, R5, A3>(
self: Stream<A, E, R>,
that: Stream<A2, E2, R2>,
s: S,
f: (
s: S,
pullLeft: Effect<A, Option<E>, R3>,
pullRight: Effect<A2, Option<E2>, R4>,
) => Effect<Exit<readonly [A3, S], Option<E2 | E>>, never, R5>,
): Stream<A3, E2 | E, R | R2 | R3 | R4 | R5>;
};combineChunks
Combines the chunks from this stream and the specified stream by repeatedly applying the function f to extract a chunk using both sides and conceptually "offer" it to the destination stream. f can maintain some internal state to control the combining process, with the initial state being specified by s.
Signature
declare const combineChunks: {
<A2, E2, R2, S, R3, E, A, R4, R5, A3>(
that: Stream<A2, E2, R2>,
s: S,
f: (
s: S,
pullLeft: Effect<Chunk<A>, Option<E>, R3>,
pullRight: Effect<Chunk<A2>, Option<E2>, R4>,
) => Effect<Exit<readonly [Chunk<A3>, S], Option<E2 | E>>, never, R5>,
): <R>(self: Stream<A, E, R>) => Stream<A3, E2 | E, R2 | R3 | R4 | R5 | R>;
<R, A2, E2, R2, S, R3, E, A, R4, R5, A3>(
self: Stream<A, E, R>,
that: Stream<A2, E2, R2>,
s: S,
f: (
s: S,
pullLeft: Effect<Chunk<A>, Option<E>, R3>,
pullRight: Effect<Chunk<A2>, Option<E2>, R4>,
) => Effect<Exit<readonly [Chunk<A3>, S], Option<E2 | E>>, never, R5>,
): Stream<A3, E2 | E, R | R2 | R3 | R4 | R5>;
};Concatenates the specified stream with this stream, resulting in a stream that emits the elements from this stream and then the elements from the specified stream.
Signature
declare const concat: {
<A2, E2, R2>(
that: Stream<A2, E2, R2>,
): <A, E, R>(self: Stream<A, E, R>) => Stream<A2 | A, E2 | E, R2 | R>;
<A, E, R, A2, E2, R2>(
self: Stream<A, E, R>,
that: Stream<A2, E2, R2>,
): Stream<A | A2, E | E2, R | R2>;
};Example
import { Effect, Stream } from "effect"
const s1 = Stream.make(1, 2, 3)
const s2 = Stream.make(4, 5)
const stream = Stream.concat(s1, s2)
Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// { _id: 'Chunk', values: [ 1, 2, 3, 4, 5 ] }Composes this stream with the specified stream to create a cartesian product of elements. The right stream would be run multiple times, for every element in the left stream.
See also Stream.zip for the more common point-wise variant.
Signature
declare const cross: {
<AR, ER, RR>(
right: Stream<AR, ER, RR>,
): <AL, EL, RL>(left: Stream<AL, EL, RL>) => Stream<[AL, AR], ER | EL, RR | RL>;
<AL, ER, RR, AR, EL, RL>(
left: Stream<AL, ER, RR>,
right: Stream<AR, EL, RL>,
): Stream<[AL, AR], ER | EL, RR | RL>;
};Example
import { Effect, Stream } from "effect"
const s1 = Stream.make(1, 2, 3)
const s2 = Stream.make("a", "b")
const product = Stream.cross(s1, s2)
Effect.runPromise(Stream.runCollect(product)).then(console.log)
// {
// _id: "Chunk",
// values: [
// [ 1, "a" ], [ 1, "b" ], [ 2, "a" ], [ 2, "b" ], [ 3, "a" ], [ 3, "b" ]
// ]
// }Composes this stream with the specified stream to create a cartesian product of elements, but keeps only elements from left stream. The right stream would be run multiple times, for every element in the left stream.
See also Stream.zipLeft for the more common point-wise variant.
Signature
declare const crossLeft: {
<AR, ER, RR>(
right: Stream<AR, ER, RR>,
): <AL, EL, RL>(left: Stream<AL, EL, RL>) => Stream<AL, ER | EL, RR | RL>;
<AL, EL, RL, AR, ER, RR>(
left: Stream<AL, EL, RL>,
right: Stream<AR, ER, RR>,
): Stream<AL, EL | ER, RL | RR>;
};crossRight
Composes this stream with the specified stream to create a cartesian product of elements, but keeps only elements from the right stream. The left stream would be run multiple times, for every element in the right stream.
See also Stream.zipRight for the more common point-wise variant.
Signature
declare const crossRight: {
<AR, ER, RR>(
right: Stream<AR, ER, RR>,
): <AL, EL, RL>(left: Stream<AL, EL, RL>) => Stream<AR, ER | EL, RR | RL>;
<AL, EL, RL, AR, ER, RR>(
left: Stream<AL, EL, RL>,
right: Stream<AR, ER, RR>,
): Stream<AR, EL | ER, RL | RR>;
};Composes this stream with the specified stream to create a cartesian product of elements with a specified function. The right stream would be run multiple times, for every element in the left stream.
See also Stream.zipWith for the more common point-wise variant.
Signature
declare const crossWith: {
<AR, ER, RR, AL, A>(
right: Stream<AR, ER, RR>,
f: (left: AL, right: AR) => A,
): <EL, RL>(left: Stream<AL, EL, RL>) => Stream<A, ER | EL, RR | RL>;
<AL, EL, RL, AR, ER, RR, A>(
left: Stream<AL, EL, RL>,
right: Stream<AR, ER, RR>,
f: (left: AL, right: AR) => A,
): Stream<A, EL | ER, RL | RR>;
};Delays the emission of values by holding new values for a set duration. If no new values arrive during that time the value is emitted, however if a new value is received during the holding period the previous value is discarded and the process is repeated with the new value.
This operator is useful if you have a stream of "bursty" events which eventually settle down and you only need the final event of the burst. For example, a search engine may only want to initiate a search after a user has paused typing so as to not prematurely recommend results.
Signature
declare const debounce: {
(duration: DurationInput): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E, R>;
<A, E, R>(self: Stream<A, E, R>, duration: DurationInput): Stream<A, E, R>;
};Example
import { Effect, Stream } from "effect"
let last = Date.now()
const log = (message: string) =>
Effect.sync(() => {
const end = Date.now()
console.log(`${message} after ${end - last}ms`)
last = end
})
const stream = Stream.make(1, 2, 3).pipe(
Stream.concat(
Stream.fromEffect(Effect.sleep("200 millis").pipe(Effect.as(4))), // Emit 4 after 200 ms
),
Stream.concat(Stream.make(5, 6)), // Continue with more rapid values
Stream.concat(
Stream.fromEffect(Effect.sleep("150 millis").pipe(Effect.as(7))), // Emit 7 after 150 ms
),
Stream.concat(Stream.make(8)),
Stream.tap((n) => log(`Received ${n}`)),
Stream.debounce("100 millis"), // Only emit values after a pause of at least 100 milliseconds,
Stream.tap((n) => log(`> Emitted ${n}`)),
)
Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// Received 1 after 5ms
// Received 2 after 2ms
// Received 3 after 0ms
// > Emitted 3 after 104ms
// Received 4 after 99ms
// Received 5 after 1ms
// Received 6 after 0ms
// > Emitted 6 after 101ms
// Received 7 after 50ms
// Received 8 after 1ms
// > Emitted 8 after 101ms
// { _id: 'Chunk', values: [ 3, 6, 8 ] }distributedWith
More powerful version of Stream.broadcast. Allows to provide a function that determines what queues should receive which elements. The decide function will receive the indices of the queues in the resulting list.
Signature
declare const distributedWith: {
<N extends number, A>(options: {
readonly decide: (a: A) => Effect.Effect<Predicate<number>>;
readonly maximumLag: number;
readonly size: N;
}): <E, R>(
self: Stream<A, E, R>,
) => Effect<TupleOf<N, Dequeue<Exit<A, Option<E>>>>, never, Scope | R>;
<A, E, R, N extends number>(
self: Stream<A, E, R>,
options: {
readonly decide: (a: A) => Effect.Effect<Predicate<number>>;
readonly maximumLag: number;
readonly size: N;
},
): Effect<TupleOf<N, Dequeue<Exit<A, Option<E>>>>, never, Scope | R>;
};distributedWithDynamic
More powerful version of Stream.distributedWith. This returns a function that will produce new queues and corresponding indices. You can also provide a function that will be executed after the final events are enqueued in all queues. Shutdown of the queues is handled by the driver. Downstream users can also shutdown queues manually. In this case the driver will continue but no longer backpressure on them.
Signature
declare const distributedWithDynamic: {
<A>(options: {
readonly decide: (a: A) => Effect.Effect<Predicate<number>, never, never>;
readonly maximumLag: number;
}): <E, R>(
self: Stream<A, E, R>,
) => Effect<Effect<[number, Dequeue<Exit<A, Option<E>>>], never, never>, never, Scope | R>;
<A, E, R>(
self: Stream<A, E, R>,
options: {
readonly decide: (a: A) => Effect.Effect<Predicate<number>, never, never>;
readonly maximumLag: number;
},
): Effect<Effect<[number, Dequeue<Exit<A, Option<E>>>], never, never>, never, Scope | R>;
};Converts this stream to a stream that executes its effects but emits no elements. Useful for sequencing effects using streams:
Signature
declare const drain: <A, E, R>(self: Stream<A, E, R>) => Stream<never, E, R>;Example
import { Effect, Stream } from "effect"
// We create a stream and immediately drain it.
const stream = Stream.range(1, 6).pipe(Stream.drain)
Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// { _id: 'Chunk', values: [] }Drains the provided stream in the background for as long as this stream is running. If this stream ends before other, other will be interrupted. If other fails, this stream will fail with that error.
Signature
declare const drainFork: {
<A2, E2, R2>(
that: Stream<A2, E2, R2>,
): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E2 | E, R2 | R>;
<A, E, R, A2, E2, R2>(self: Stream<A, E, R>, that: Stream<A2, E2, R2>): Stream<A, E | E2, R | R2>;
};Drops the specified number of elements from this stream.
Signature
declare const drop: {
(n: number): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E, R>;
<A, E, R>(self: Stream<A, E, R>, n: number): Stream<A, E, R>;
};Drops the last specified number of elements from this stream.
Signature
declare const dropRight: {
(n: number): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E, R>;
<A, E, R>(self: Stream<A, E, R>, n: number): Stream<A, E, R>;
};Drops all elements of the stream until the specified predicate evaluates to true.
Signature
declare const dropUntil: {
<A>(predicate: Predicate<NoInfer<A>>): <E, R>(self: Stream<A, E, R>) => Stream<A, E, R>;
<A, E, R>(self: Stream<A, E, R>, predicate: Predicate<A>): Stream<A, E, R>;
};dropUntilEffect
Drops all elements of the stream until the specified effectful predicate evaluates to true.
Signature
declare const dropUntilEffect: {
<A, E2, R2>(
predicate: (a: NoInfer<A>) => Effect<boolean, E2, R2>,
): <E, R>(self: Stream<A, E, R>) => Stream<A, E2 | E, R2 | R>;
<A, E, R, E2, R2>(
self: Stream<A, E, R>,
predicate: (a: NoInfer<A>) => Effect<boolean, E2, R2>,
): Stream<A, E | E2, R | R2>;
};Drops all elements of the stream for as long as the specified predicate evaluates to true.
Signature
declare const dropWhile: {
<A>(predicate: Predicate<NoInfer<A>>): <E, R>(self: Stream<A, E, R>) => Stream<A, E, R>;
<A, E, R>(self: Stream<A, E, R>, predicate: Predicate<A>): Stream<A, E, R>;
};dropWhileEffect
Drops all elements of the stream for as long as the specified predicate produces an effect that evalutates to true
Signature
declare const dropWhileEffect: {
<A, E2, R2>(
predicate: (a: NoInfer<A>) => Effect<boolean, E2, R2>,
): <E, R>(self: Stream<A, E, R>) => Stream<A, E2 | E, R2 | R>;
<A, E, R, E2, R2>(
self: Stream<A, E, R>,
predicate: (a: A) => Effect<boolean, E2, R2>,
): Stream<A, E | E2, R | R2>;
};Returns a stream whose failures and successes have been lifted into an Either. The resulting stream cannot fail, because the failures have been exposed as part of the Either success case.
Signature
declare const either: <A, E, R>(self: Stream<A, E, R>) => Stream<Either.Either<A, E>, never, R>;Executes the provided finalizer after this stream's finalizers run.
Signature
declare const ensuring: {
<X, R2>(
finalizer: Effect<X, never, R2>,
): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E, R2 | R>;
<A, E, R, X, R2>(self: Stream<A, E, R>, finalizer: Effect<X, never, R2>): Stream<A, E, R | R2>;
};Example
import { Console, Effect, Stream } from "effect"
const program = Stream.fromEffect(Console.log("Application Logic.")).pipe(
Stream.concat(Stream.finalizer(Console.log("Finalizing the stream"))),
Stream.ensuring(Console.log("Doing some other works after stream's finalization")),
)
Effect.runPromise(Stream.runCollect(program)).then(console.log)
// Application Logic.
// Finalizing the stream
// Doing some other works after stream's finalization
// { _id: 'Chunk', values: [ undefined, undefined ] }ensuringWith
Executes the provided finalizer after this stream's finalizers run.
Signature
declare const ensuringWith: {
<E, R2>(
finalizer: (exit: Exit<unknown, E>) => Effect<unknown, never, R2>,
): <A, R>(self: Stream<A, E, R>) => Stream<A, E, R2 | R>;
<A, E, R, R2>(
self: Stream<A, E, R>,
finalizer: (exit: Exit<unknown, E>) => Effect<unknown, never, R2>,
): Stream<A, E, R | R2>;
};Performs a filter and map in a single step.
Signature
declare const filterMap: {
<A, B>(pf: (a: A) => Option<B>): <E, R>(self: Stream<A, E, R>) => Stream<B, E, R>;
<A, E, R, B>(self: Stream<A, E, R>, pf: (a: A) => Option<B>): Stream<B, E, R>;
};filterMapEffect
Performs an effectful filter and map in a single step.
Signature
declare const filterMapEffect: {
<A, A2, E2, R2>(
pf: (a: A) => Option<Effect<A2, E2, R2>>,
): <E, R>(self: Stream<A, E, R>) => Stream<A2, E2 | E, R2 | R>;
<A, E, R, A2, E2, R2>(
self: Stream<A, E, R>,
pf: (a: A) => Option<Effect<A2, E2, R2>>,
): Stream<A2, E | E2, R | R2>;
};filterMapWhile
Transforms all elements of the stream for as long as the specified partial function is defined.
Signature
declare const filterMapWhile: {
<A, A2>(pf: (a: A) => Option<A2>): <E, R>(self: Stream<A, E, R>) => Stream<A2, E, R>;
<A, E, R, A2>(self: Stream<A, E, R>, pf: (a: A) => Option<A2>): Stream<A2, E, R>;
};filterMapWhileEffect
Effectfully transforms all elements of the stream for as long as the specified partial function is defined.
Signature
declare const filterMapWhileEffect: {
<A, A2, E2, R2>(
pf: (a: A) => Option<Effect<A2, E2, R2>>,
): <E, R>(self: Stream<A, E, R>) => Stream<A2, E2 | E, R2 | R>;
<A, E, R, A2, E2, R2>(
self: Stream<A, E, R>,
pf: (a: A) => Option<Effect<A2, E2, R2>>,
): Stream<A2, E | E2, R | R2>;
};Repeats this stream forever.
Signature
declare const forever: <A, E, R>(self: Stream<A, E, R>) => Stream<A, E, R>;Specialized version of haltWhen which halts the evaluation of this stream after the given duration.
An element in the process of being pulled will not be interrupted when the given duration completes. See interruptAfter for this behavior.
Signature
declare const haltAfter: {
(duration: DurationInput): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E, R>;
<A, E, R>(self: Stream<A, E, R>, duration: DurationInput): Stream<A, E, R>;
};Halts the evaluation of this stream when the provided effect completes. The given effect will be forked as part of the returned stream, and its success will be discarded.
An element in the process of being pulled will not be interrupted when the effect completes. See interruptWhen for this behavior.
If the effect completes with a failure, the stream will emit that failure.
Signature
declare const haltWhen: {
<X, E2, R2>(
effect: Effect<X, E2, R2>,
): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E2 | E, R2 | R>;
<A, E, R, X, E2, R2>(self: Stream<A, E, R>, effect: Effect<X, E2, R2>): Stream<A, E | E2, R | R2>;
};haltWhenDeferred
Halts the evaluation of this stream when the provided promise resolves.
If the promise completes with a failure, the stream will emit that failure.
Signature
declare const haltWhenDeferred: {
<X, E2>(deferred: Deferred<X, E2>): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E2 | E, R>;
<A, E, R, X, E2>(self: Stream<A, E, R>, deferred: Deferred<X, E2>): Stream<A, E | E2, R>;
};The identity pipeline, which does not modify streams in any way.
Signature
declare const identity: <A, E = never, R = never>() => Stream<A, E, R>;interleave
Interleaves this stream and the specified stream deterministically by alternating pulling values from this stream and the specified stream. When one stream is exhausted all remaining values in the other stream will be pulled.
Signature
declare const interleave: {
<A2, E2, R2>(
that: Stream<A2, E2, R2>,
): <A, E, R>(self: Stream<A, E, R>) => Stream<A2 | A, E2 | E, R2 | R>;
<A, E, R, A2, E2, R2>(
self: Stream<A, E, R>,
that: Stream<A2, E2, R2>,
): Stream<A | A2, E | E2, R | R2>;
};Example
import { Effect, Stream } from "effect"
const s1 = Stream.make(1, 2, 3)
const s2 = Stream.make(4, 5, 6)
const stream = Stream.interleave(s1, s2)
Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// { _id: 'Chunk', values: [ 1, 4, 2, 5, 3, 6 ] }interleaveWith
Combines this stream and the specified stream deterministically using the stream of boolean values pull to control which stream to pull from next. A value of true indicates to pull from this stream and a value of false indicates to pull from the specified stream. Only consumes as many elements as requested by the pull stream. If either this stream or the specified stream are exhausted further requests for values from that stream will be ignored.
Signature
declare const interleaveWith: {
<A2, E2, R2, E3, R3>(
that: Stream<A2, E2, R2>,
decider: Stream<boolean, E3, R3>,
): <A, E, R>(self: Stream<A, E, R>) => Stream<A2 | A, E2 | E3 | E, R2 | R3 | R>;
<A, E, R, A2, E2, R2, E3, R3>(
self: Stream<A, E, R>,
that: Stream<A2, E2, R2>,
decider: Stream<boolean, E3, R3>,
): Stream<A | A2, E | E2 | E3, R | R2 | R3>;
};Example
import { Effect, Stream } from "effect"
const s1 = Stream.make(1, 3, 5, 7, 9)
const s2 = Stream.make(2, 4, 6, 8, 10)
const booleanStream = Stream.make(true, false, false).pipe(Stream.forever)
const stream = Stream.interleaveWith(s1, s2, booleanStream)
Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// {
// _id: 'Chunk',
// values: [
// 1, 2, 4, 3, 6,
// 8, 5, 10, 7, 9
// ]
// }interruptAfter
Specialized version of Stream.interruptWhen which interrupts the evaluation of this stream after the given Duration.
Signature
declare const interruptAfter: {
(duration: DurationInput): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E, R>;
<A, E, R>(self: Stream<A, E, R>, duration: DurationInput): Stream<A, E, R>;
};interruptWhen
Interrupts the evaluation of this stream when the provided effect completes. The given effect will be forked as part of this stream, and its success will be discarded. This combinator will also interrupt any in-progress element being pulled from upstream.
If the effect completes with a failure before the stream completes, the returned stream will emit that failure.
Signature
declare const interruptWhen: {
<X, E2, R2>(
effect: Effect<X, E2, R2>,
): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E2 | E, R2 | R>;
<A, E, R, X, E2, R2>(self: Stream<A, E, R>, effect: Effect<X, E2, R2>): Stream<A, E | E2, R | R2>;
};interruptWhenDeferred
Interrupts the evaluation of this stream when the provided promise resolves. This combinator will also interrupt any in-progress element being pulled from upstream.
If the promise completes with a failure, the stream will emit that failure.
Signature
declare const interruptWhenDeferred: {
<X, E2>(deferred: Deferred<X, E2>): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E2 | E, R>;
<A, E, R, X, E2>(self: Stream<A, E, R>, deferred: Deferred<X, E2>): Stream<A, E | E2, R>;
};intersperse
Intersperse stream with provided element.
Signature
declare const intersperse: {
<A2>(element: A2): <A, E, R>(self: Stream<A, E, R>) => Stream<A2 | A, E, R>;
<A, E, R, A2>(self: Stream<A, E, R>, element: A2): Stream<A | A2, E, R>;
};Example
import { Effect, Stream } from "effect"
const stream = Stream.make(1, 2, 3, 4, 5).pipe(Stream.intersperse(0))
Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// {
// _id: 'Chunk',
// values: [
// 1, 0, 2, 0, 3,
// 0, 4, 0, 5
// ]
// }intersperseAffixes
Intersperse the specified element, also adding a prefix and a suffix.
Signature
declare const intersperseAffixes: {
<A2, A3, A4>(options: {
readonly end: A4;
readonly middle: A3;
readonly start: A2;
}): <A, E, R>(self: Stream<A, E, R>) => Stream<A2 | A3 | A4 | A, E, R>;
<A, E, R, A2, A3, A4>(
self: Stream<A, E, R>,
options: {
readonly end: A4;
readonly middle: A3;
readonly start: A2;
},
): Stream<A | A2 | A3 | A4, E, R>;
};Example
import { Effect, Stream } from "effect"
const stream = Stream.make(1, 2, 3, 4, 5).pipe(
Stream.intersperseAffixes({
start: "[",
middle: "-",
end: "]",
}),
)
Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// {
// _id: 'Chunk',
// values: [
// '[', 1, '-', 2, '-',
// 3, '-', 4, '-', 5,
// ']'
// ]
// }Returns a stream whose failure and success channels have been mapped by the specified onFailure and onSuccess functions.
Signature
declare const mapBoth: {
<E, E2, A, A2>(options: {
readonly onFailure: (e: E) => E2;
readonly onSuccess: (a: A) => A2;
}): <R>(self: Stream<A, E, R>) => Stream<A2, E2, R>;
<A, E, R, E2, A2>(
self: Stream<A, E, R>,
options: {
readonly onFailure: (e: E) => E2;
readonly onSuccess: (a: A) => A2;
},
): Stream<A2, E2, R>;
};Merges this stream and the specified stream together.
New produced stream will terminate when both specified stream terminate if no termination strategy is specified.
Signature
declare const merge: {
<A2, E2, R2>(
that: Stream<A2, E2, R2>,
options?: {
readonly haltStrategy?: HaltStrategy.HaltStrategyInput;
},
): <A, E, R>(self: Stream<A, E, R>) => Stream<A2 | A, E2 | E, R2 | R>;
<A, E, R, A2, E2, R2>(
self: Stream<A, E, R>,
that: Stream<A2, E2, R2>,
options?: {
readonly haltStrategy?: HaltStrategy.HaltStrategyInput;
},
): Stream<A | A2, E | E2, R | R2>;
};Example
import { Effect, Schedule, Stream } from "effect"
const s1 = Stream.make(1, 2, 3).pipe(Stream.schedule(Schedule.spaced("100 millis")))
const s2 = Stream.make(4, 5, 6).pipe(Stream.schedule(Schedule.spaced("200 millis")))
const stream = Stream.merge(s1, s2)
Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// { _id: 'Chunk', values: [ 1, 4, 2, 3, 5, 6 ] }Merges a variable list of streams in a non-deterministic fashion. Up to n streams may be consumed in parallel and up to outputBuffer chunks may be buffered by this operator.
Signature
declare const mergeAll: {
(options: {
readonly bufferSize?: number;
readonly concurrency: number | "unbounded";
}): <A, E, R>(streams: Iterable<Stream<A, E, R>>) => Stream<A, E, R>;
<A, E, R>(
streams: Iterable<Stream<A, E, R>>,
options: {
readonly bufferSize?: number;
readonly concurrency: number | "unbounded";
},
): Stream<A, E, R>;
};mergeEither
Merges this stream and the specified stream together to produce a stream of eithers.
Signature
declare const mergeEither: {
<A2, E2, R2>(
that: Stream<A2, E2, R2>,
): <A, E, R>(self: Stream<A, E, R>) => Stream<Either<A2, A>, E2 | E, R2 | R>;
<A, E, R, A2, E2, R2>(
self: Stream<A, E, R>,
that: Stream<A2, E2, R2>,
): Stream<Either<A2, A>, E | E2, R | R2>;
};Merges this stream and the specified stream together, discarding the values from the right stream.
Signature
declare const mergeLeft: {
<AR, ER, RR>(
right: Stream<AR, ER, RR>,
): <AL, EL, RL>(left: Stream<AL, EL, RL>) => Stream<AL, ER | EL, RR | RL>;
<AL, EL, RL, AR, ER, RR>(
left: Stream<AL, EL, RL>,
right: Stream<AR, ER, RR>,
): Stream<AL, EL | ER, RL | RR>;
};mergeRight
Merges this stream and the specified stream together, discarding the values from the left stream.
Signature
declare const mergeRight: {
<AR, ER, RR>(
right: Stream<AR, ER, RR>,
): <AL, EL, RL>(left: Stream<AL, EL, RL>) => Stream<AR, ER | EL, RR | RL>;
<AL, EL, RL, AR, ER, RR>(
left: Stream<AL, EL, RL>,
right: Stream<AR, ER, RR>,
): Stream<AR, EL | ER, RL | RR>;
};Merges this stream and the specified stream together to a common element type with the specified mapping functions.
New produced stream will terminate when both specified stream terminate if no termination strategy is specified.
Signature
declare const mergeWith: {
<A2, E2, R2, A, A3, A4>(
other: Stream<A2, E2, R2>,
options: {
readonly haltStrategy?: HaltStrategy.HaltStrategyInput;
readonly onOther: (a2: A2) => A4;
readonly onSelf: (a: A) => A3;
},
): <E, R>(self: Stream<A, E, R>) => Stream<A3 | A4, E2 | E, R2 | R>;
<A, E, R, A2, E2, R2, A3, A4>(
self: Stream<A, E, R>,
other: Stream<A2, E2, R2>,
options: {
readonly haltStrategy?: HaltStrategy.HaltStrategyInput;
readonly onOther: (a2: A2) => A4;
readonly onSelf: (a: A) => A3;
},
): Stream<A3 | A4, E | E2, R | R2>;
};Example
import { Effect, Schedule, Stream } from "effect"
const s1 = Stream.make("1", "2", "3").pipe(Stream.schedule(Schedule.spaced("100 millis")))
const s2 = Stream.make(4.1, 5.3, 6.2).pipe(Stream.schedule(Schedule.spaced("200 millis")))
const stream = Stream.mergeWith(s1, s2, {
onSelf: (s) => parseInt(s),
onOther: (n) => Math.floor(n),
})
Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// { _id: 'Chunk', values: [ 1, 4, 2, 3, 5, 6 ] }Returns a combined string resulting from concatenating each of the values from the stream.
Signature
declare const mkString: <E, R>(self: Stream<string, E, R>) => Effect.Effect<string, E, R>;Runs the specified effect if this stream ends.
Signature
declare const onDone: {
<X, R2>(
cleanup: () => Effect<X, never, R2>,
): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E, R2 | R>;
<A, E, R, X, R2>(
self: Stream<A, E, R>,
cleanup: () => Effect<X, never, R2>,
): Stream<A, E, R | R2>;
};Runs the specified effect if this stream fails, providing the error to the effect if it exists.
Note: Unlike Effect.onError there is no guarantee that the provided effect will not be interrupted.
Signature
declare const onError: {
<E, X, R2>(
cleanup: (cause: Cause<E>) => Effect<X, never, R2>,
): <A, R>(self: Stream<A, E, R>) => Stream<A, E, R2 | R>;
<A, E, R, X, R2>(
self: Stream<A, E, R>,
cleanup: (cause: Cause<E>) => Effect<X, never, R2>,
): Stream<A, E, R | R2>;
};Splits a stream into two substreams based on a predicate.
Details
The Stream.partition function splits a stream into two parts: one for elements that satisfy the predicate (evaluated to true) and another for those that do not (evaluated to false).
The faster stream may advance up to bufferSize elements ahead of the slower one.
See
partitionEitherfor partitioning a stream based on effectful conditions.
Signature
declare const partition: {
<C, B, A = C>(
refinement: Refinement<NoInfer<A>, B>,
options?: {
bufferSize?: number;
},
): <E, R>(
self: Stream<C, E, R>,
) => Effect<
[excluded: Stream<Exclude<C, B>, E, never>, satisfying: Stream<B, E, never>],
E,
Scope | R
>;
<A>(
predicate: Predicate<A>,
options?: {
bufferSize?: number;
},
): <E, R>(
self: Stream<A, E, R>,
) => Effect<[excluded: Stream<A, E, never>, satisfying: Stream<A, E, never>], E, Scope | R>;
<C, E, R, B, A = C>(
self: Stream<C, E, R>,
refinement: Refinement<A, B>,
options?: {
bufferSize?: number;
},
): Effect<
[excluded: Stream<Exclude<C, B>, E, never>, satisfying: Stream<B, E, never>],
E,
Scope | R
>;
<A, E, R>(
self: Stream<A, E, R>,
predicate: Predicate<A>,
options?: {
bufferSize?: number;
},
): Effect<[excluded: Stream<A, E, never>, satisfying: Stream<A, E, never>], E, Scope | R>;
};Example
(Partitioning a Stream into Even and Odd Numbers)
import { Effect, Stream } from "effect"
const partition = Stream.range(1, 9).pipe(Stream.partition((n) => n % 2 === 0, { bufferSize: 5 }))
const program = Effect.scoped(
Effect.gen(function* () {
const [odds, evens] = yield* partition
console.log(yield* Stream.runCollect(odds))
console.log(yield* Stream.runCollect(evens))
}),
)
Effect.runPromise(program)
// { _id: 'Chunk', values: [ 1, 3, 5, 7, 9 ] }
// { _id: 'Chunk', values: [ 2, 4, 6, 8 ] }partitionEither
Splits a stream into two substreams based on an effectful condition.
Details
The Stream.partitionEither function is used to divide a stream into two parts: one for elements that satisfy a condition producing Either.left values, and another for those that produce Either.right values. This function applies an effectful predicate to each element in the stream to determine which substream it belongs to.
The faster stream may advance up to bufferSize elements ahead of the slower one.
See
partitionfor partitioning a stream based on simple conditions.
Signature
declare const partitionEither: {
<A, A3, A2, E2, R2>(
predicate: (a: NoInfer<A>) => Effect<Either<A3, A2>, E2, R2>,
options?: {
readonly bufferSize?: number;
},
): <E, R>(
self: Stream<A, E, R>,
) => Effect<
[left: Stream<A2, E2 | E, never>, right: Stream<A3, E2 | E, never>],
E2 | E,
Scope | R2 | R
>;
<A, E, R, A3, A2, E2, R2>(
self: Stream<A, E, R>,
predicate: (a: A) => Effect<Either<A3, A2>, E2, R2>,
options?: {
readonly bufferSize?: number;
},
): Effect<
[left: Stream<A2, E | E2, never>, right: Stream<A3, E | E2, never>],
E | E2,
Scope | R | R2
>;
};Example
(Partitioning a Stream with an Effectful Predicate)
import { Effect, Either, Stream } from "effect"
const partition = Stream.range(1, 9).pipe(
Stream.partitionEither((n) => Effect.succeed(n % 2 === 0 ? Either.right(n) : Either.left(n)), {
bufferSize: 5,
}),
)
const program = Effect.scoped(
Effect.gen(function* () {
const [evens, odds] = yield* partition
console.log(yield* Stream.runCollect(evens))
console.log(yield* Stream.runCollect(odds))
}),
)
Effect.runPromise(program)
// { _id: 'Chunk', values: [ 1, 3, 5, 7, 9 ] }
// { _id: 'Chunk', values: [ 2, 4, 6, 8 ] }Peels off enough material from the stream to construct a Z using the provided Sink and then returns both the Z and the rest of the Stream in a scope. Like all scoped values, the provided stream is valid only within the scope.
Signature
declare const peel: {
<A2, A, E2, R2>(
sink: Sink<A2, A, A, E2, R2>,
): <E, R>(self: Stream<A, E, R>) => Effect<[A2, Stream<A, E, never>], E2 | E, Scope | R2 | R>;
<A, E, R, A2, E2, R2>(
self: Stream<A, E, R>,
sink: Sink<A2, A, A, E2, R2>,
): Effect<[A2, Stream<A, E, never>], E | E2, Scope | R | R2>;
};pipeThrough
Pipes all of the values from this stream through the provided sink.
See also Stream.transduce.
Signature
declare const pipeThrough: {
<A2, A, L, E2, R2>(
sink: Sink<A2, A, L, E2, R2>,
): <E, R>(self: Stream<A, E, R>) => Stream<L, E2 | E, R2 | R>;
<A, E, R, A2, L, E2, R2>(
self: Stream<A, E, R>,
sink: Sink<A2, A, L, E2, R2>,
): Stream<L, E | E2, R | R2>;
};pipeThroughChannel
Pipes all the values from this stream through the provided channel.
Signature
declare const pipeThroughChannel: {
<R2, E, E2, A, A2>(
channel: Channel<Chunk<A2>, Chunk<A>, E2, E, unknown, unknown, R2>,
): <R>(self: Stream<A, E, R>) => Stream<A2, E2, R2 | R>;
<R, R2, E, E2, A, A2>(
self: Stream<A, E, R>,
channel: Channel<Chunk<A2>, Chunk<A>, E2, E, unknown, unknown, R2>,
): Stream<A2, E2, R | R2>;
};pipeThroughChannelOrFail
Pipes all values from this stream through the provided channel, passing through any error emitted by this stream unchanged.
Signature
declare const pipeThroughChannelOrFail: {
<R2, E, E2, A, A2>(
chan: Channel<Chunk<A2>, Chunk<A>, E2, E, unknown, unknown, R2>,
): <R>(self: Stream<A, E, R>) => Stream<A2, E | E2, R2 | R>;
<R, R2, E, E2, A, A2>(
self: Stream<A, E, R>,
chan: Channel<Chunk<A2>, Chunk<A>, E2, E, unknown, unknown, R2>,
): Stream<A2, E | E2, R | R2>;
};Emits the provided chunk before emitting any other value.
Signature
declare const prepend: {
<B>(values: Chunk<B>): <A, E, R>(self: Stream<A, E, R>) => Stream<B | A, E, R>;
<A, E, R, B>(self: Stream<A, E, R>, values: Chunk<B>): Stream<A | B, E, R>;
};Re-chunks the elements of the stream into chunks of n elements each. The last chunk might contain less than n elements.
Signature
declare const rechunk: {
(n: number): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E, R>;
<A, E, R>(self: Stream<A, E, R>, n: number): Stream<A, E, R>;
};Repeats the entire stream using the specified schedule. The stream will execute normally, and then repeat again according to the provided schedule.
Signature
declare const repeat: {
<B, R2>(
schedule: Schedule<B, unknown, R2>,
): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E, R2 | R>;
<A, E, R, B, R2>(self: Stream<A, E, R>, schedule: Schedule<B, unknown, R2>): Stream<A, E, R | R2>;
};Example
import { Effect, Schedule, Stream } from "effect"
const stream = Stream.repeat(Stream.succeed(1), Schedule.forever)
Effect.runPromise(Stream.runCollect(stream.pipe(Stream.take(5)))).then(console.log)
// { _id: 'Chunk', values: [ 1, 1, 1, 1, 1 ] }repeatEither
Repeats the entire stream using the specified schedule. The stream will execute normally, and then repeat again according to the provided schedule. The schedule output will be emitted at the end of each repetition.
Signature
declare const repeatEither: {
<B, R2>(
schedule: Schedule<B, unknown, R2>,
): <A, E, R>(self: Stream<A, E, R>) => Stream<Either<A, B>, E, R2 | R>;
<A, E, R, B, R2>(
self: Stream<A, E, R>,
schedule: Schedule<B, unknown, R2>,
): Stream<Either<A, B>, E, R | R2>;
};repeatElements
Repeats each element of the stream using the provided schedule. Repetitions are done in addition to the first execution, which means using Schedule.recurs(1) actually results in the original effect, plus an additional recurrence, for a total of two repetitions of each value in the stream.
Signature
declare const repeatElements: {
<B, R2>(
schedule: Schedule<B, unknown, R2>,
): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E, R2 | R>;
<A, E, R, B, R2>(self: Stream<A, E, R>, schedule: Schedule<B, unknown, R2>): Stream<A, E, R | R2>;
};repeatElementsWith
Repeats each element of the stream using the provided schedule. When the schedule is finished, then the output of the schedule will be emitted into the stream. Repetitions are done in addition to the first execution, which means using Schedule.recurs(1) actually results in the original effect, plus an additional recurrence, for a total of two repetitions of each value in the stream.
This function accepts two conversion functions, which allow the output of this stream and the output of the provided schedule to be unified into a single type. For example, Either or similar data type.
Signature
declare const repeatElementsWith: {
<B, R2, A, C>(
schedule: Schedule<B, unknown, R2>,
options: {
readonly onElement: (a: A) => C;
readonly onSchedule: (b: B) => C;
},
): <E, R>(self: Stream<A, E, R>) => Stream<C, E, R2 | R>;
<A, E, R, B, R2, C>(
self: Stream<A, E, R>,
schedule: Schedule<B, unknown, R2>,
options: {
readonly onElement: (a: A) => C;
readonly onSchedule: (b: B) => C;
},
): Stream<C, E, R | R2>;
};repeatWith
Repeats the entire stream using the specified schedule. The stream will execute normally, and then repeat again according to the provided schedule. The schedule output will be emitted at the end of each repetition and can be unified with the stream elements using the provided functions.
Signature
declare const repeatWith: {
<B, R2, A, C>(
schedule: Schedule<B, unknown, R2>,
options: {
readonly onElement: (a: A) => C;
readonly onSchedule: (b: B) => C;
},
): <E, R>(self: Stream<A, E, R>) => Stream<C, E, R2 | R>;
<A, E, R, B, R2, C>(
self: Stream<A, E, R>,
schedule: Schedule<B, unknown, R2>,
options: {
readonly onElement: (a: A) => C;
readonly onSchedule: (b: B) => C;
},
): Stream<C, E, R | R2>;
};When the stream fails, retry it according to the given schedule
This retries the entire stream, so will re-execute all of the stream's acquire operations.
The schedule is reset as soon as the first element passes through the stream again.
Signature
declare const retry: {
<E, R2, X>(
policy: Schedule<X, NoInfer<E>, R2>,
): <A, R>(self: Stream<A, E, R>) => Stream<A, E, R2 | R>;
<A, E, R, X, R2>(
self: Stream<A, E, R>,
policy: Schedule<X, NoInfer<E>, R2>,
): Stream<A, E, R | R2>;
};Statefully maps over the elements of this stream to produce all intermediate results of type S given an initial S.
Signature
declare const scan: {
<S, A>(s: S, f: (s: S, a: A) => S): <E, R>(self: Stream<A, E, R>) => Stream<S, E, R>;
<A, E, R, S>(self: Stream<A, E, R>, s: S, f: (s: S, a: A) => S): Stream<S, E, R>;
};Example
import { Effect, Stream } from "effect"
const stream = Stream.range(1, 6).pipe(Stream.scan(0, (a, b) => a + b))
Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// { _id: 'Chunk', values: [ 0, 1, 3, 6, 10, 15, 21 ] }scanEffect
Statefully and effectfully maps over the elements of this stream to produce all intermediate results of type S given an initial S.
Signature
declare const scanEffect: {
<S, A, E2, R2>(
s: S,
f: (s: S, a: A) => Effect<S, E2, R2>,
): <E, R>(self: Stream<A, E, R>) => Stream<S, E2 | E, R2 | R>;
<A, E, R, S, E2, R2>(
self: Stream<A, E, R>,
s: S,
f: (s: S, a: A) => Effect<S, E2, R2>,
): Stream<S, E | E2, R | R2>;
};scanReduce
Statefully maps over the elements of this stream to produce all intermediate results.
See also Stream.scan.
Signature
declare const scanReduce: {
<A2, A>(f: (a2: A2 | A, a: A) => A2): <E, R>(self: Stream<A, E, R>) => Stream<A2 | A, E, R>;
<A, E, R, A2>(self: Stream<A, E, R>, f: (a2: A | A2, a: A) => A2): Stream<A | A2, E, R>;
};scanReduceEffect
Statefully and effectfully maps over the elements of this stream to produce all intermediate results.
See also Stream.scanEffect.
Signature
declare const scanReduceEffect: {
<A2, A, E2, R2>(
f: (a2: A2 | A, a: A) => Effect<A2 | A, E2, R2>,
): <E, R>(self: Stream<A, E, R>) => Stream<A2 | A, E2 | E, R2 | R>;
<A, E, R, A2, E2, R2>(
self: Stream<A, E, R>,
f: (a2: A | A2, a: A) => Effect<A | A2, E2, R2>,
): Stream<A | A2, E | E2, R | R2>;
};Schedules the output of the stream using the provided schedule.
Signature
declare const schedule: {
<X, A0, R2, A>(
schedule: Schedule<X, A0, R2>,
): <E, R>(self: Stream<A, E, R>) => Stream<A, E, R2 | R>;
<A, E, R, X, A0, R2>(self: Stream<A, E, R>, schedule: Schedule<X, A0, R2>): Stream<A, E, R | R2>;
};scheduleWith
Schedules the output of the stream using the provided schedule and emits its output at the end (if schedule is finite). Uses the provided function to align the stream and schedule outputs on the same type.
Signature
declare const scheduleWith: {
<B, A0, R2, A, C>(
schedule: Schedule<B, A0, R2>,
options: {
readonly onElement: (a: A) => C;
readonly onSchedule: (b: B) => C;
},
): <E, R>(self: Stream<A, E, R>) => Stream<C, E, R2 | R>;
<A, E, R, B, A0, R2, C>(
self: Stream<A, E, R>,
schedule: Schedule<B, A0, R2>,
options: {
readonly onElement: (a: A) => C;
readonly onSchedule: (b: B) => C;
},
): Stream<C, E, R | R2>;
};Emits a sliding window of n elements.
Signature
declare const sliding: {
(chunkSize: number): <A, E, R>(self: Stream<A, E, R>) => Stream<Chunk<A>, E, R>;
<A, E, R>(self: Stream<A, E, R>, chunkSize: number): Stream<Chunk<A>, E, R>;
};Example
import { pipe, Stream } from "effect"
pipe(Stream.make(1, 2, 3, 4), Stream.sliding(2), Stream.runCollect)
// => Chunk(Chunk(1, 2), Chunk(2, 3), Chunk(3, 4))slidingSize
Like sliding, but with a configurable stepSize parameter.
Signature
declare const slidingSize: {
(chunkSize: number, stepSize: number): <A, E, R>(self: Stream<A, E, R>) => Stream<Chunk<A>, E, R>;
<A, E, R>(self: Stream<A, E, R>, chunkSize: number, stepSize: number): Stream<Chunk<A>, E, R>;
};Converts an option on values into an option on errors.
Signature
declare const some: <A, E, R>(
self: Stream<Option.Option<A>, E, R>,
) => Stream<A, Option.Option<E>, R>;someOrElse
Extracts the optional value, or returns the given 'default'.
Signature
declare const someOrElse: {
<A2>(fallback: LazyArg<A2>): <A, E, R>(self: Stream<Option<A>, E, R>) => Stream<A2 | A, E, R>;
<A, E, R, A2>(self: Stream<Option<A>, E, R>, fallback: LazyArg<A2>): Stream<A | A2, E, R>;
};someOrFail
Extracts the optional value, or fails with the given error 'e'.
Signature
declare const someOrFail: {
<E2>(error: LazyArg<E2>): <A, E, R>(self: Stream<Option<A>, E, R>) => Stream<A, E2 | E, R>;
<A, E, R, E2>(self: Stream<Option<A>, E, R>, error: LazyArg<E2>): Stream<A, E | E2, R>;
};Splits elements based on a predicate or refinement.
Signature
declare const split: {
<A, B>(
refinement: Refinement<NoInfer<A>, B>,
): <E, R>(self: Stream<A, E, R>) => Stream<Chunk<Exclude<A, B>>, E, R>;
<A>(predicate: Predicate<NoInfer<A>>): <E, R>(self: Stream<A, E, R>) => Stream<Chunk<A>, E, R>;
<A, E, R, B>(
self: Stream<A, E, R>,
refinement: Refinement<A, B>,
): Stream<Chunk<Exclude<A, B>>, E, R>;
<A, E, R>(self: Stream<A, E, R>, predicate: Predicate<A>): Stream<Chunk<A>, E, R>;
};Example
import { pipe, Stream } from "effect"
pipe(
Stream.range(1, 10),
Stream.split((n) => n % 4 === 0),
Stream.runCollect,
)
// => Chunk(Chunk(1, 2, 3), Chunk(5, 6, 7), Chunk(9))splitOnChunk
Splits elements on a delimiter and transforms the splits into desired output.
Signature
declare const splitOnChunk: {
<A>(delimiter: Chunk<A>): <E, R>(self: Stream<A, E, R>) => Stream<Chunk<A>, E, R>;
<A, E, R>(self: Stream<A, E, R>, delimiter: Chunk<A>): Stream<Chunk<A>, E, R>;
};Takes the specified number of elements from this stream.
Signature
declare const take: {
(n: number): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E, R>;
<A, E, R>(self: Stream<A, E, R>, n: number): Stream<A, E, R>;
};Example
import { Effect, Stream } from "effect"
const stream = Stream.take(
Stream.iterate(0, (n) => n + 1),
5,
)
Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// { _id: 'Chunk', values: [ 0, 1, 2, 3, 4 ] }Takes the last specified number of elements from this stream.
Signature
declare const takeRight: {
(n: number): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E, R>;
<A, E, R>(self: Stream<A, E, R>, n: number): Stream<A, E, R>;
};Example
import { Effect, Stream } from "effect"
const stream = Stream.takeRight(Stream.make(1, 2, 3, 4, 5, 6), 3)
Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// { _id: 'Chunk', values: [ 4, 5, 6 ] }Takes all elements of the stream until the specified predicate evaluates to true.
Signature
declare const takeUntil: {
<A>(predicate: Predicate<NoInfer<A>>): <E, R>(self: Stream<A, E, R>) => Stream<A, E, R>;
<A, E, R>(self: Stream<A, E, R>, predicate: Predicate<A>): Stream<A, E, R>;
};Example
import { Effect, Stream } from "effect"
const stream = Stream.takeUntil(
Stream.iterate(0, (n) => n + 1),
(n) => n === 4,
)
Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// { _id: 'Chunk', values: [ 0, 1, 2, 3, 4 ] }takeUntilEffect
Takes all elements of the stream until the specified effectual predicate evaluates to true.
Signature
declare const takeUntilEffect: {
<A, E2, R2>(
predicate: (a: NoInfer<A>) => Effect<boolean, E2, R2>,
): <E, R>(self: Stream<A, E, R>) => Stream<A, E2 | E, R2 | R>;
<A, E, R, E2, R2>(
self: Stream<A, E, R>,
predicate: (a: A) => Effect<boolean, E2, R2>,
): Stream<A, E | E2, R | R2>;
};Takes all elements of the stream for as long as the specified predicate evaluates to true.
Signature
declare const takeWhile: {
<A, B>(refinement: Refinement<NoInfer<A>, B>): <E, R>(self: Stream<A, E, R>) => Stream<B, E, R>;
<A>(predicate: Predicate<NoInfer<A>>): <E, R>(self: Stream<A, E, R>) => Stream<A, E, R>;
<A, E, R, B>(self: Stream<A, E, R>, refinement: Refinement<A, B>): Stream<B, E, R>;
<A, E, R>(self: Stream<A, E, R>, predicate: Predicate<A>): Stream<A, E, R>;
};Example
import { Effect, Stream } from "effect"
const stream = Stream.takeWhile(
Stream.iterate(0, (n) => n + 1),
(n) => n < 5,
)
Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// { _id: 'Chunk', values: [ 0, 1, 2, 3, 4 ] }tapErrorCause
Returns a stream that effectfully "peeks" at the cause of failure of the stream.
Signature
declare const tapErrorCause: {
<E, X, E2, R2>(
f: (cause: Cause<NoInfer<E>>) => Effect<X, E2, R2>,
): <A, R>(self: Stream<A, E, R>) => Stream<A, E | E2, R2 | R>;
<A, E, R, X, E2, R2>(
self: Stream<A, E, R>,
f: (cause: Cause<E>) => Effect<X, E2, R2>,
): Stream<A, E | E2, R | R2>;
};Delays the chunks of this stream according to the given bandwidth parameters using the token bucket algorithm. Allows for burst in the processing of elements by allowing the token bucket to accumulate tokens up to a units + burst threshold. The weight of each chunk is determined by the cost function.
If using the "enforce" strategy, chunks that do not meet the bandwidth constraints are dropped. If using the "shape" strategy, chunks are delayed until they can be emitted without exceeding the bandwidth constraints.
Defaults to the "shape" strategy.
Signature
declare const throttle: {
<A>(options: {
readonly burst?: number;
readonly cost: (chunk: Chunk.Chunk<A>) => number;
readonly duration: Duration.DurationInput;
readonly strategy?: "enforce" | "shape";
readonly units: number;
}): <E, R>(self: Stream<A, E, R>) => Stream<A, E, R>;
<A, E, R>(
self: Stream<A, E, R>,
options: {
readonly burst?: number;
readonly cost: (chunk: Chunk.Chunk<A>) => number;
readonly duration: Duration.DurationInput;
readonly strategy?: "enforce" | "shape";
readonly units: number;
},
): Stream<A, E, R>;
};Example
import { Chunk, Effect, Schedule, Stream } from "effect"
let last = Date.now()
const log = (message: string) =>
Effect.sync(() => {
const end = Date.now()
console.log(`${message} after ${end - last}ms`)
last = end
})
const stream = Stream.fromSchedule(Schedule.spaced("50 millis")).pipe(
Stream.take(6),
Stream.tap((n) => log(`Received ${n}`)),
Stream.throttle({
cost: Chunk.size,
duration: "100 millis",
units: 1,
}),
Stream.tap((n) => log(`> Emitted ${n}`)),
)
Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// Received 0 after 56ms
// > Emitted 0 after 0ms
// Received 1 after 52ms
// > Emitted 1 after 48ms
// Received 2 after 52ms
// > Emitted 2 after 49ms
// Received 3 after 52ms
// > Emitted 3 after 48ms
// Received 4 after 52ms
// > Emitted 4 after 47ms
// Received 5 after 52ms
// > Emitted 5 after 49ms
// { _id: 'Chunk', values: [ 0, 1, 2, 3, 4, 5 ] }throttleEffect
Delays the chunks of this stream according to the given bandwidth parameters using the token bucket algorithm. Allows for burst in the processing of elements by allowing the token bucket to accumulate tokens up to a units + burst threshold. The weight of each chunk is determined by the effectful costFn function.
If using the "enforce" strategy, chunks that do not meet the bandwidth constraints are dropped. If using the "shape" strategy, chunks are delayed until they can be emitted without exceeding the bandwidth constraints.
Defaults to the "shape" strategy.
Signature
declare const throttleEffect: {
<A, E2, R2>(options: {
readonly burst?: number;
readonly cost: (chunk: Chunk.Chunk<A>) => Effect.Effect<number, E2, R2>;
readonly duration: Duration.DurationInput;
readonly strategy?: "enforce" | "shape";
readonly units: number;
}): <E, R>(self: Stream<A, E, R>) => Stream<A, E2 | E, R2 | R>;
<A, E, R, E2, R2>(
self: Stream<A, E, R>,
options: {
readonly burst?: number;
readonly cost: (chunk: Chunk.Chunk<A>) => Effect.Effect<number, E2, R2>;
readonly duration: Duration.DurationInput;
readonly strategy?: "enforce" | "shape";
readonly units: number;
},
): Stream<A, E | E2, R | R2>;
};Ends the stream if it does not produce a value after the specified duration.
Signature
declare const timeout: {
(duration: DurationInput): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E, R>;
<A, E, R>(self: Stream<A, E, R>, duration: DurationInput): Stream<A, E, R>;
};timeoutFail
Fails the stream with given error if it does not produce a value after d duration.
Signature
declare const timeoutFail: {
<E2>(
error: LazyArg<E2>,
duration: DurationInput,
): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E2 | E, R>;
<A, E, R, E2>(
self: Stream<A, E, R>,
error: LazyArg<E2>,
duration: DurationInput,
): Stream<A, E | E2, R>;
};timeoutFailCause
Fails the stream with given cause if it does not produce a value after d duration.
Signature
declare const timeoutFailCause: {
<E2>(
cause: LazyArg<Cause<E2>>,
duration: DurationInput,
): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E2 | E, R>;
<A, E, R, E2>(
self: Stream<A, E, R>,
cause: LazyArg<Cause<E2>>,
duration: DurationInput,
): Stream<A, E | E2, R>;
};Switches the stream if it does not produce a value after the specified duration.
Signature
declare const timeoutTo: {
<A2, E2, R2>(
duration: DurationInput,
that: Stream<A2, E2, R2>,
): <A, E, R>(self: Stream<A, E, R>) => Stream<A2 | A, E2 | E, R2 | R>;
<A, E, R, A2, E2, R2>(
self: Stream<A, E, R>,
duration: DurationInput,
that: Stream<A2, E2, R2>,
): Stream<A | A2, E | E2, R | R2>;
};Applies the transducer to the stream and emits its outputs.
Signature
declare const transduce: {
<A2, A, E2, R2>(
sink: Sink<A2, A, A, E2, R2>,
): <E, R>(self: Stream<A, E, R>) => Stream<A2, E2 | E, R2 | R>;
<A, E, R, A2, E2, R2>(
self: Stream<A, E, R>,
sink: Sink<A2, A, A, E2, R2>,
): Stream<A2, E | E2, R | R2>;
};Returns the specified stream if the given condition is satisfied, otherwise returns an empty stream.
Signature
declare const when: {
(test: LazyArg<boolean>): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E, R>;
<A, E, R>(self: Stream<A, E, R>, test: LazyArg<boolean>): Stream<A, E, R>;
};whenCaseEffect
Returns the stream when the given partial function is defined for the given effectful value, otherwise returns an empty stream.
Signature
declare const whenCaseEffect: {
<A, A2, E2, R2>(
pf: (a: A) => Option<Stream<A2, E2, R2>>,
): <E, R>(self: Effect<A, E, R>) => Stream<A2, E2 | E, R2 | R>;
<A, E, R, A2, E2, R2>(
self: Effect<A, E, R>,
pf: (a: A) => Option<Stream<A2, E2, R2>>,
): Stream<A2, E | E2, R | R2>;
};whenEffect
Returns the stream if the given effectful condition is satisfied, otherwise returns an empty stream.
Signature
declare const whenEffect: {
<E2, R2>(
effect: Effect<boolean, E2, R2>,
): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E2 | E, R2 | R>;
<A, E, R, E2, R2>(
self: Stream<A, E, R>,
effect: Effect<boolean, E2, R2>,
): Stream<A, E | E2, R | R2>;
};Zipping
Zips this stream with another point-wise and emits tuples of elements from both streams.
The new stream will end when one of the sides ends.
Signature
declare const zip: {
<A2, E2, R2>(
that: Stream<A2, E2, R2>,
): <A, E, R>(self: Stream<A, E, R>) => Stream<[A, A2], E2 | E, R2 | R>;
<A, E, R, A2, E2, R2>(
self: Stream<A, E, R>,
that: Stream<A2, E2, R2>,
): Stream<[A, A2], E | E2, R | R2>;
};Example
import { Effect, Stream } from "effect"
// We create two streams and zip them together.
const stream = Stream.zip(Stream.make(1, 2, 3, 4, 5, 6), Stream.make("a", "b", "c"))
Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// { _id: 'Chunk', values: [ [ 1, 'a' ], [ 2, 'b' ], [ 3, 'c' ] ] }Zips this stream with another point-wise, creating a new stream of pairs of elements from both sides.
The defaults defaultLeft and defaultRight will be used if the streams have different lengths and one of the streams has ended before the other.
Signature
declare const zipAll: {
<A2, E2, R2, A>(options: {
readonly defaultOther: A2;
readonly defaultSelf: A;
readonly other: Stream<A2, E2, R2>;
}): <E, R>(self: Stream<A, E, R>) => Stream<[A, A2], E2 | E, R2 | R>;
<A, E, R, A2, E2, R2>(
self: Stream<A, E, R>,
options: {
readonly defaultOther: A2;
readonly defaultSelf: A;
readonly other: Stream<A2, E2, R2>;
},
): Stream<[A, A2], E | E2, R | R2>;
};Example
import { Effect, Stream } from "effect"
const stream = Stream.zipAll(Stream.make(1, 2, 3, 4, 5, 6), {
other: Stream.make("a", "b", "c"),
defaultSelf: 0,
defaultOther: "x",
})
Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// { _id: "Chunk", values: [ [ 1, "a" ], [ 2, "b" ], [ 3, "c" ], [ 4, "x" ], [ 5, "x" ], [ 6, "x" ] ] }zipAllLeft
Zips this stream with another point-wise, and keeps only elements from this stream.
The provided default value will be used if the other stream ends before this one.
Signature
declare const zipAllLeft: {
<A2, E2, R2, A>(
that: Stream<A2, E2, R2>,
defaultLeft: A,
): <E, R>(self: Stream<A, E, R>) => Stream<A, E2 | E, R2 | R>;
<A, E, R, A2, E2, R2>(
self: Stream<A, E, R>,
that: Stream<A2, E2, R2>,
defaultLeft: A,
): Stream<A, E | E2, R | R2>;
};zipAllRight
Zips this stream with another point-wise, and keeps only elements from the other stream.
The provided default value will be used if this stream ends before the other one.
Signature
declare const zipAllRight: {
<A2, E2, R2>(
that: Stream<A2, E2, R2>,
defaultRight: A2,
): <A, E, R>(self: Stream<A, E, R>) => Stream<A2, E2 | E, R2 | R>;
<A, E, R, A2, E2, R2>(
self: Stream<A, E, R>,
that: Stream<A2, E2, R2>,
defaultRight: A2,
): Stream<A2, E | E2, R | R2>;
};zipAllSortedByKey
Zips this stream that is sorted by distinct keys and the specified stream that is sorted by distinct keys to produce a new stream that is sorted by distinct keys. Combines values associated with each key into a tuple, using the specified values defaultLeft and defaultRight to fill in missing values.
This allows zipping potentially unbounded streams of data by key in constant space but the caller is responsible for ensuring that the streams are sorted by distinct keys.
Signature
declare const zipAllSortedByKey: {
<A2, E2, R2, A, K>(options: {
readonly defaultOther: A2;
readonly defaultSelf: A;
readonly order: Order.Order<K>;
readonly other: Stream<readonly [K, A2], E2, R2>;
}): <E, R>(self: Stream<readonly [K, A], E, R>) => Stream<[K, [A, A2]], E2 | E, R2 | R>;
<K, A, E, R, A2, E2, R2>(
self: Stream<readonly [K, A], E, R>,
options: {
readonly defaultOther: A2;
readonly defaultSelf: A;
readonly order: Order.Order<K>;
readonly other: Stream<readonly [K, A2], E2, R2>;
},
): Stream<[K, [A, A2]], E | E2, R | R2>;
};zipAllSortedByKeyLeft
Zips this stream that is sorted by distinct keys and the specified stream that is sorted by distinct keys to produce a new stream that is sorted by distinct keys. Keeps only values from this stream, using the specified value default to fill in missing values.
This allows zipping potentially unbounded streams of data by key in constant space but the caller is responsible for ensuring that the streams are sorted by distinct keys.
Signature
declare const zipAllSortedByKeyLeft: {
<A2, E2, R2, A, K>(options: {
readonly defaultSelf: A;
readonly order: Order.Order<K>;
readonly other: Stream<readonly [K, A2], E2, R2>;
}): <E, R>(self: Stream<readonly [K, A], E, R>) => Stream<[K, A], E2 | E, R2 | R>;
<K, A, E, R, A2, E2, R2>(
self: Stream<readonly [K, A], E, R>,
options: {
readonly defaultSelf: A;
readonly order: Order.Order<K>;
readonly other: Stream<readonly [K, A2], E2, R2>;
},
): Stream<[K, A], E | E2, R | R2>;
};zipAllSortedByKeyRight
Zips this stream that is sorted by distinct keys and the specified stream that is sorted by distinct keys to produce a new stream that is sorted by distinct keys. Keeps only values from that stream, using the specified value default to fill in missing values.
This allows zipping potentially unbounded streams of data by key in constant space but the caller is responsible for ensuring that the streams are sorted by distinct keys.
Signature
declare const zipAllSortedByKeyRight: {
<K, A2, E2, R2>(options: {
readonly defaultOther: A2;
readonly order: Order.Order<K>;
readonly other: Stream<readonly [K, A2], E2, R2>;
}): <A, E, R>(self: Stream<readonly [K, A], E, R>) => Stream<[K, A2], E2 | E, R2 | R>;
<A, E, R, K, A2, E2, R2>(
self: Stream<readonly [K, A], E, R>,
options: {
readonly defaultOther: A2;
readonly order: Order.Order<K>;
readonly other: Stream<readonly [K, A2], E2, R2>;
},
): Stream<[K, A2], E | E2, R | R2>;
};zipAllSortedByKeyWith
Zips this stream that is sorted by distinct keys and the specified stream that is sorted by distinct keys to produce a new stream that is sorted by distinct keys. Uses the functions left, right, and both to handle the cases where a key and value exist in this stream, that stream, or both streams.
This allows zipping potentially unbounded streams of data by key in constant space but the caller is responsible for ensuring that the streams are sorted by distinct keys.
Signature
declare const zipAllSortedByKeyWith: {
<K, A2, E2, R2, A, A3>(options: {
readonly onBoth: (a: A, a2: A2) => A3;
readonly onOther: (a2: A2) => A3;
readonly onSelf: (a: A) => A3;
readonly order: Order.Order<K>;
readonly other: Stream<readonly [K, A2], E2, R2>;
}): <E, R>(self: Stream<readonly [K, A], E, R>) => Stream<[K, A3], E2 | E, R2 | R>;
<K, A, E, R, A2, E2, R2, A3>(
self: Stream<readonly [K, A], E, R>,
options: {
readonly onBoth: (a: A, a2: A2) => A3;
readonly onOther: (a2: A2) => A3;
readonly onSelf: (a: A) => A3;
readonly order: Order.Order<K>;
readonly other: Stream<readonly [K, A2], E2, R2>;
},
): Stream<[K, A3], E | E2, R | R2>;
};zipAllWith
Zips this stream with another point-wise. The provided functions will be used to create elements for the composed stream.
The functions left and right will be used if the streams have different lengths and one of the streams has ended before the other.
Signature
declare const zipAllWith: {
<A2, E2, R2, A, A3>(options: {
readonly onBoth: (a: A, a2: A2) => A3;
readonly onOther: (a2: A2) => A3;
readonly onSelf: (a: A) => A3;
readonly other: Stream<A2, E2, R2>;
}): <E, R>(self: Stream<A, E, R>) => Stream<A3, E2 | E, R2 | R>;
<A, E, R, A2, E2, R2, A3>(
self: Stream<A, E, R>,
options: {
readonly onBoth: (a: A, a2: A2) => A3;
readonly onOther: (a2: A2) => A3;
readonly onSelf: (a: A) => A3;
readonly other: Stream<A2, E2, R2>;
},
): Stream<A3, E | E2, R | R2>;
};Example
import { Effect, Stream } from "effect"
const stream = Stream.zipAllWith(Stream.make(1, 2, 3, 4, 5, 6), {
other: Stream.make("a", "b", "c"),
onSelf: (n) => [n, "x"],
onOther: (s) => [0, s],
onBoth: (n, s) => [n - s.length, s],
})
Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// { _id: "Chunk", values: [ [ 0, "a" ], [ 1, "b" ], [ 2, "c" ], [ 4, "x" ], [ 5, "x" ], [ 6, "x" ] ] }zipFlatten
Zips this stream with another point-wise and emits tuples of elements from both streams.
The new stream will end when one of the sides ends.
Signature
declare const zipFlatten: {
<A2, E2, R2>(that: Stream<A2, E2, R2>): <A extends readonly Array<any>, E, R>(self: Stream<A, E, R>) => Stream<[...Array<A>, A2], E2 | E, R2 | R>;
<A extends readonly Array<any>, E, R, A2, E2, R2>(self: Stream<A, E, R>, that: Stream<A2, E2, R2>): Stream<[...Array<A>, A2], E | E2, R | R2>;
}Zips the two streams so that when a value is emitted by either of the two streams, it is combined with the latest value from the other stream to produce a result.
Note: tracking the latest value is done on a per-chunk basis. That means that emitted elements that are not the last value in chunks will never be used for zipping.
Signature
declare const zipLatest: {
<AR, ER, RR>(
right: Stream<AR, ER, RR>,
): <AL, EL, RL>(left: Stream<AL, EL, RL>) => Stream<[AL, AR], ER | EL, RR | RL>;
<AL, EL, RL, AR, ER, RR>(
left: Stream<AL, EL, RL>,
right: Stream<AR, ER, RR>,
): Stream<[AL, AR], EL | ER, RL | RR>;
};Example
import { Effect, Schedule, Stream } from "effect"
const s1 = Stream.make(1, 2, 3).pipe(Stream.schedule(Schedule.spaced("1 second")))
const s2 = Stream.make("a", "b", "c", "d").pipe(Stream.schedule(Schedule.spaced("500 millis")))
const stream = Stream.zipLatest(s1, s2)
Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// { _id: "Chunk", values: [ [ 1, "a" ], [ 1, "b" ], [ 2, "b" ], [ 2, "c" ], [ 2, "d" ], [ 3, "d" ] ] }zipLatestAll
Zips multiple streams so that when a value is emitted by any of the streams, it is combined with the latest values from the other streams to produce a result.
Note: tracking the latest value is done on a per-chunk basis. That means that emitted elements that are not the last value in chunks will never be used for zipping.
Signature
declare const zipLatestAll: <T extends ReadonlyArray<Stream<any, any, any>>>(
...streams: T
) => Stream<
[T[number]] extends [never]
? never
: { [K in keyof T]: T[K] extends Stream<infer A, infer _E, infer _R> ? A : never },
[T[number]] extends [never]
? never
: T[number] extends Stream<infer _A, infer _E, infer _R>
? _E
: never,
[T[number]] extends [never]
? never
: T[number] extends Stream<infer _A, infer _E, infer _R>
? _R
: never
>;Example
import { Stream, Schedule, Console, Effect } from "effect"
const stream = Stream.zipLatestAll(
Stream.fromSchedule(Schedule.spaced("1 millis")),
Stream.fromSchedule(Schedule.spaced("2 millis")),
Stream.fromSchedule(Schedule.spaced("4 millis")),
).pipe(Stream.take(6), Stream.tap(Console.log))
Effect.runPromise(Stream.runDrain(stream))
// Output:
// [ 0, 0, 0 ]
// [ 1, 0, 0 ]
// [ 1, 1, 0 ]
// [ 2, 1, 0 ]
// [ 3, 1, 0 ]
// [ 3, 1, 1 ]
// .....zipLatestWith
Zips the two streams so that when a value is emitted by either of the two streams, it is combined with the latest value from the other stream to produce a result.
Note: tracking the latest value is done on a per-chunk basis. That means that emitted elements that are not the last value in chunks will never be used for zipping.
Signature
declare const zipLatestWith: {
<AR, ER, RR, AL, A>(
right: Stream<AR, ER, RR>,
f: (left: AL, right: AR) => A,
): <EL, RL>(left: Stream<AL, EL, RL>) => Stream<A, ER | EL, RR | RL>;
<AL, EL, RL, AR, ER, RR, A>(
left: Stream<AL, EL, RL>,
right: Stream<AR, ER, RR>,
f: (left: AL, right: AR) => A,
): Stream<A, EL | ER, RL | RR>;
};Zips this stream with another point-wise, but keeps only the outputs of left stream.
The new stream will end when one of the sides ends.
Signature
declare const zipLeft: {
<AR, ER, RR>(
right: Stream<AR, ER, RR>,
): <AL, EL, RL>(left: Stream<AL, EL, RL>) => Stream<AL, ER | EL, RR | RL>;
<AL, EL, RL, AR, ER, RR>(
left: Stream<AL, EL, RL>,
right: Stream<AR, ER, RR>,
): Stream<AL, EL | ER, RL | RR>;
};Zips this stream with another point-wise, but keeps only the outputs of the right stream.
The new stream will end when one of the sides ends.
Signature
declare const zipRight: {
<AR, ER, RR>(
right: Stream<AR, ER, RR>,
): <AL, EL, RL>(left: Stream<AL, EL, RL>) => Stream<AR, ER | EL, RR | RL>;
<AL, EL, RL, AR, ER, RR>(
left: Stream<AL, EL, RL>,
right: Stream<AR, ER, RR>,
): Stream<AR, EL | ER, RL | RR>;
};Zips this stream with another point-wise and applies the function to the paired elements.
The new stream will end when one of the sides ends.
Signature
declare const zipWith: {
<AR, ER, RR, AL, A>(
right: Stream<AR, ER, RR>,
f: (left: AL, right: AR) => A,
): <EL, RL>(left: Stream<AL, EL, RL>) => Stream<A, ER | EL, RR | RL>;
<AL, EL, RL, AR, ER, RR, A>(
left: Stream<AL, EL, RL>,
right: Stream<AR, ER, RR>,
f: (left: AL, right: AR) => A,
): Stream<A, EL | ER, RL | RR>;
};Example
import { Effect, Stream } from "effect"
// We create two streams and zip them with custom logic.
const stream = Stream.zipWith(Stream.make(1, 2, 3, 4, 5, 6), Stream.make("a", "b", "c"), (n, s) => [
n - s.length,
s,
])
Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// { _id: 'Chunk', values: [ [ 0, 'a' ], [ 1, 'b' ], [ 2, 'c' ] ] }zipWithChunks
Zips this stream with another point-wise and applies the function to the paired elements.
The new stream will end when one of the sides ends.
Signature
declare const zipWithChunks: {
<A2, E2, R2, A, A3>(
that: Stream<A2, E2, R2>,
f: (left: Chunk<A>, right: Chunk<A2>) => readonly [Chunk<A3>, Either<Chunk<A2>, Chunk<A>>],
): <E, R>(self: Stream<A, E, R>) => Stream<A3, E2 | E, R2 | R>;
<A, E, R, A2, E2, R2, A3>(
self: Stream<A, E, R>,
that: Stream<A2, E2, R2>,
f: (left: Chunk<A>, right: Chunk<A2>) => readonly [Chunk<A3>, Either<Chunk<A2>, Chunk<A>>],
): Stream<A3, E | E2, R | R2>;
};zipWithIndex
Zips this stream together with the index of elements.
Signature
declare const zipWithIndex: <A, E, R>(self: Stream<A, E, R>) => Stream<[A, number], E, R>;Example
import { Effect, Stream } from "effect"
const stream = Stream.make("Mary", "James", "Robert", "Patricia")
const indexedStream = Stream.zipWithIndex(stream)
Effect.runPromise(Stream.runCollect(indexedStream)).then(console.log)
// {
// _id: 'Chunk',
// values: [ [ 'Mary', 0 ], [ 'James', 1 ], [ 'Robert', 2 ], [ 'Patricia', 3 ] ]
// }zipWithNext
Zips each element with the next element if present.
Signature
declare const zipWithNext: <A, E, R>(self: Stream<A, E, R>) => Stream<[A, Option.Option<A>], E, R>;Example
import { Chunk, Effect, Stream } from "effect"
const stream = Stream.zipWithNext(Stream.make(1, 2, 3, 4))
Effect.runPromise(Stream.runCollect(stream)).then((chunk) => console.log(Chunk.toArray(chunk)))
// [
// [ 1, { _id: 'Option', _tag: 'Some', value: 2 } ],
// [ 2, { _id: 'Option', _tag: 'Some', value: 3 } ],
// [ 3, { _id: 'Option', _tag: 'Some', value: 4 } ],
// [ 4, { _id: 'Option', _tag: 'None' } ]
// ]zipWithPrevious
Zips each element with the previous element. Initially accompanied by None.
Signature
declare const zipWithPrevious: <A, E, R>(
self: Stream<A, E, R>,
) => Stream<[Option.Option<A>, A], E, R>;Example
import { Chunk, Effect, Stream } from "effect"
const stream = Stream.zipWithPrevious(Stream.make(1, 2, 3, 4))
Effect.runPromise(Stream.runCollect(stream)).then((chunk) => console.log(Chunk.toArray(chunk)))
// [
// [ { _id: 'Option', _tag: 'None' }, 1 ],
// [ { _id: 'Option', _tag: 'Some', value: 1 }, 2 ],
// [ { _id: 'Option', _tag: 'Some', value: 2 }, 3 ],
// [ { _id: 'Option', _tag: 'Some', value: 3 }, 4 ]
// ]zipWithPreviousAndNext
Zips each element with both the previous and next element.
Signature
declare const zipWithPreviousAndNext: <A, E, R>(
self: Stream<A, E, R>,
) => Stream<[Option.Option<A>, A, Option.Option<A>], E, R>;Example
import { Chunk, Effect, Stream } from "effect"
const stream = Stream.zipWithPreviousAndNext(Stream.make(1, 2, 3, 4))
Effect.runPromise(Stream.runCollect(stream)).then((chunk) => console.log(Chunk.toArray(chunk)))
// [
// [
// { _id: 'Option', _tag: 'None' },
// 1,
// { _id: 'Option', _tag: 'Some', value: 2 }
// ],
// [
// { _id: 'Option', _tag: 'Some', value: 1 },
// 2,
// { _id: 'Option', _tag: 'Some', value: 3 }
// ],
// [
// { _id: 'Option', _tag: 'Some', value: 2 },
// 3,
// { _id: 'Option', _tag: 'Some', value: 4 }
// ],
// [
// { _id: 'Option', _tag: 'Some', value: 3 },
// 4,
// { _id: 'Option', _tag: 'None' }
// ]
// ]
Merges a struct of streams into a single stream of tagged values.