(
options:
| {
readonly capacity: "unbounded"
readonly replay?: number | undefined
readonly idleTimeToLive?: Duration.Input | undefined
}
| {
readonly capacity: number
readonly strategy?: "sliding" | "dropping" | "suspend" | undefined
readonly replay?: number | undefined
readonly idleTimeToLive?: Duration.Input | undefined
}
): <A, E, R>(
self: Stream<A, E, R>
) => Effect.Effect<Stream<A, E>, never, Scope.Scope | R>
<A, E, R>(
self: Stream<A, E, R>,
options:
| {
readonly capacity: "unbounded"
readonly replay?: number | undefined
readonly idleTimeToLive?: Duration.Input | undefined
}
| {
readonly capacity: number
readonly strategy?: "sliding" | "dropping" | "suspend" | undefined
readonly replay?: number | undefined
readonly idleTimeToLive?: Duration.Input | undefined
}
): Effect.Effect<Stream<A, E>, never, Scope.Scope | R>Returns a new Stream that multicasts the original stream, subscribing when the first consumer starts.
Details
The upstream continues running while there is at least one consumer and is finalized after the last one exits.
If idleTimeToLive is set, the upstream is kept alive for that duration so a later subscriber can continue from
the next element instead of restarting.
Example (Sharing a stream)
import { Console, Effect, Stream } from "effect"
Effect.runPromise(
Effect.scoped(
Effect.gen(function*() {
const shared = yield* Stream.make(1, 2, 3).pipe(
Stream.share({ capacity: 16 })
)
const first = yield* shared.pipe(Stream.take(1), Stream.runCollect)
const second = yield* shared.pipe(Stream.take(1), Stream.runCollect)
yield* Console.log([first, second])
})
)
)
// output: [[1], [1]]export const const share: {
(
options:
| {
readonly capacity: "unbounded"
readonly replay?: number | undefined
readonly idleTimeToLive?:
| Duration.Input
| undefined
}
| {
readonly capacity: number
readonly strategy?:
| "sliding"
| "dropping"
| "suspend"
| undefined
readonly replay?: number | undefined
readonly idleTimeToLive?:
| Duration.Input
| undefined
}
): <A, E, R>(
self: Stream<A, E, R>
) => Effect.Effect<
Stream<A, E>,
never,
Scope.Scope | R
>
<A, E, R>(
self: Stream<A, E, R>,
options:
| {
readonly capacity: "unbounded"
readonly replay?: number | undefined
readonly idleTimeToLive?:
| Duration.Input
| undefined
}
| {
readonly capacity: number
readonly strategy?:
| "sliding"
| "dropping"
| "suspend"
| undefined
readonly replay?: number | undefined
readonly idleTimeToLive?:
| Duration.Input
| undefined
}
): Effect.Effect<
Stream<A, E>,
never,
Scope.Scope | R
>
}
Returns a new Stream that multicasts the original stream, subscribing when the first consumer starts.
Details
The upstream continues running while there is at least one consumer and is finalized after the last one exits.
If idleTimeToLive is set, the upstream is kept alive for that duration so a later subscriber can continue from
the next element instead of restarting.
Example (Sharing a stream)
import { Console, Effect, Stream } from "effect"
Effect.runPromise(
Effect.scoped(
Effect.gen(function*() {
const shared = yield* Stream.make(1, 2, 3).pipe(
Stream.share({ capacity: 16 })
)
const first = yield* shared.pipe(Stream.take(1), Stream.runCollect)
const second = yield* shared.pipe(Stream.take(1), Stream.runCollect)
yield* Console.log([first, second])
})
)
)
// output: [[1], [1]]
share: {
(
options: | {
readonly capacity: "unbounded"
readonly replay?: number | undefined
readonly idleTimeToLive?:
| Duration.Input
| undefined
}
| {
readonly capacity: number
readonly strategy?:
| "sliding"
| "dropping"
| "suspend"
| undefined
readonly replay?: number | undefined
readonly idleTimeToLive?:
| Duration.Input
| undefined
}
options: {
readonly capacity: "unbounded"capacity: "unbounded"
readonly replay?: number | undefinedreplay?: number | undefined
readonly idleTimeToLive?: Duration.Input | undefinedidleTimeToLive?: import DurationDuration.type Duration.Input = /*unresolved*/ anyInput | undefined
} | {
readonly capacity: numbercapacity: number
readonly strategy?: | "suspend"
| "sliding"
| "dropping"
| undefined
strategy?: "sliding" | "dropping" | "suspend" | undefined
readonly replay?: number | undefinedreplay?: number | undefined
readonly idleTimeToLive?: Duration.Input | undefinedidleTimeToLive?: import DurationDuration.type Duration.Input = /*unresolved*/ anyInput | undefined
}
): <function (type parameter) A in <A, E, R>(self: Stream<A, E, R>): Effect.Effect<Stream<A, E>, never, Scope.Scope | R>A, function (type parameter) E in <A, E, R>(self: Stream<A, E, R>): Effect.Effect<Stream<A, E>, never, Scope.Scope | R>E, function (type parameter) R in <A, E, R>(self: Stream<A, E, R>): Effect.Effect<Stream<A, E>, never, Scope.Scope | R>R>(self: Stream<A, E, R>(parameter) self: {
channel: Channel.Channel<Arr.NonEmptyReadonlyArray<A>, E, void, unknown, unknown, unknown, R>;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
self: interface Stream<out A, out E = never, out R = never>A Stream<A, E, R> describes a program that can emit many A values, fail
with E, and require R.
Details
Streams are pull-based with backpressure and emit chunks to amortize effect
evaluation. They support monadic composition and error handling similar to
Effect, adapted for multiple values.
Example (Creating and consuming streams)
import { Console, Effect, Stream } from "effect"
const program = Effect.gen(function*() {
yield* Stream.make(1, 2, 3).pipe(
Stream.map((n) => n * 2),
Stream.runForEach((n) => Console.log(n))
)
})
Effect.runPromise(program)
// Output:
// 2
// 4
// 6
Stream<function (type parameter) A in <A, E, R>(self: Stream<A, E, R>): Effect.Effect<Stream<A, E>, never, Scope.Scope | R>A, function (type parameter) E in <A, E, R>(self: Stream<A, E, R>): Effect.Effect<Stream<A, E>, never, Scope.Scope | R>E, function (type parameter) R in <A, E, R>(self: Stream<A, E, R>): Effect.Effect<Stream<A, E>, never, Scope.Scope | R>R>) => import EffectEffect.type Effect.Effect = /*unresolved*/ anyEffect<interface Stream<out A, out E = never, out R = never>A Stream<A, E, R> describes a program that can emit many A values, fail
with E, and require R.
Details
Streams are pull-based with backpressure and emit chunks to amortize effect
evaluation. They support monadic composition and error handling similar to
Effect, adapted for multiple values.
Example (Creating and consuming streams)
import { Console, Effect, Stream } from "effect"
const program = Effect.gen(function*() {
yield* Stream.make(1, 2, 3).pipe(
Stream.map((n) => n * 2),
Stream.runForEach((n) => Console.log(n))
)
})
Effect.runPromise(program)
// Output:
// 2
// 4
// 6
Stream<function (type parameter) A in <A, E, R>(self: Stream<A, E, R>): Effect.Effect<Stream<A, E>, never, Scope.Scope | R>A, function (type parameter) E in <A, E, R>(self: Stream<A, E, R>): Effect.Effect<Stream<A, E>, never, Scope.Scope | R>E>, never, import ScopeScope.type Scope.Scope = /*unresolved*/ anyScope | function (type parameter) R in <A, E, R>(self: Stream<A, E, R>): Effect.Effect<Stream<A, E>, never, Scope.Scope | R>R>
<function (type parameter) A in <A, E, R>(self: Stream<A, E, R>, options: {
readonly capacity: "unbounded";
readonly replay?: number | undefined;
readonly idleTimeToLive?: Duration.Input | undefined;
} | {
readonly capacity: number;
readonly strategy?: "sliding" | "dropping" | "suspend" | undefined;
readonly replay?: number | undefined;
readonly idleTimeToLive?: Duration.Input | undefined;
}): Effect.Effect<Stream<A, E>, never, Scope.Scope | R>
A, function (type parameter) E in <A, E, R>(self: Stream<A, E, R>, options: {
readonly capacity: "unbounded";
readonly replay?: number | undefined;
readonly idleTimeToLive?: Duration.Input | undefined;
} | {
readonly capacity: number;
readonly strategy?: "sliding" | "dropping" | "suspend" | undefined;
readonly replay?: number | undefined;
readonly idleTimeToLive?: Duration.Input | undefined;
}): Effect.Effect<Stream<A, E>, never, Scope.Scope | R>
E, function (type parameter) R in <A, E, R>(self: Stream<A, E, R>, options: {
readonly capacity: "unbounded";
readonly replay?: number | undefined;
readonly idleTimeToLive?: Duration.Input | undefined;
} | {
readonly capacity: number;
readonly strategy?: "sliding" | "dropping" | "suspend" | undefined;
readonly replay?: number | undefined;
readonly idleTimeToLive?: Duration.Input | undefined;
}): Effect.Effect<Stream<A, E>, never, Scope.Scope | R>
R>(
self: Stream<A, E, R>(parameter) self: {
channel: Channel.Channel<Arr.NonEmptyReadonlyArray<A>, E, void, unknown, unknown, unknown, R>;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
self: interface Stream<out A, out E = never, out R = never>A Stream<A, E, R> describes a program that can emit many A values, fail
with E, and require R.
Details
Streams are pull-based with backpressure and emit chunks to amortize effect
evaluation. They support monadic composition and error handling similar to
Effect, adapted for multiple values.
Example (Creating and consuming streams)
import { Console, Effect, Stream } from "effect"
const program = Effect.gen(function*() {
yield* Stream.make(1, 2, 3).pipe(
Stream.map((n) => n * 2),
Stream.runForEach((n) => Console.log(n))
)
})
Effect.runPromise(program)
// Output:
// 2
// 4
// 6
Stream<function (type parameter) A in <A, E, R>(self: Stream<A, E, R>, options: {
readonly capacity: "unbounded";
readonly replay?: number | undefined;
readonly idleTimeToLive?: Duration.Input | undefined;
} | {
readonly capacity: number;
readonly strategy?: "sliding" | "dropping" | "suspend" | undefined;
readonly replay?: number | undefined;
readonly idleTimeToLive?: Duration.Input | undefined;
}): Effect.Effect<Stream<A, E>, never, Scope.Scope | R>
A, function (type parameter) E in <A, E, R>(self: Stream<A, E, R>, options: {
readonly capacity: "unbounded";
readonly replay?: number | undefined;
readonly idleTimeToLive?: Duration.Input | undefined;
} | {
readonly capacity: number;
readonly strategy?: "sliding" | "dropping" | "suspend" | undefined;
readonly replay?: number | undefined;
readonly idleTimeToLive?: Duration.Input | undefined;
}): Effect.Effect<Stream<A, E>, never, Scope.Scope | R>
E, function (type parameter) R in <A, E, R>(self: Stream<A, E, R>, options: {
readonly capacity: "unbounded";
readonly replay?: number | undefined;
readonly idleTimeToLive?: Duration.Input | undefined;
} | {
readonly capacity: number;
readonly strategy?: "sliding" | "dropping" | "suspend" | undefined;
readonly replay?: number | undefined;
readonly idleTimeToLive?: Duration.Input | undefined;
}): Effect.Effect<Stream<A, E>, never, Scope.Scope | R>
R>,
options: | {
readonly capacity: "unbounded"
readonly replay?: number | undefined
readonly idleTimeToLive?:
| Duration.Input
| undefined
}
| {
readonly capacity: number
readonly strategy?:
| "sliding"
| "dropping"
| "suspend"
| undefined
readonly replay?: number | undefined
readonly idleTimeToLive?:
| Duration.Input
| undefined
}
options: {
readonly capacity: "unbounded"capacity: "unbounded"
readonly replay?: number | undefinedreplay?: number | undefined
readonly idleTimeToLive?: Duration.Input | undefinedidleTimeToLive?: import DurationDuration.type Duration.Input = /*unresolved*/ anyInput | undefined
} | {
readonly capacity: numbercapacity: number
readonly strategy?: | "suspend"
| "sliding"
| "dropping"
| undefined
strategy?: "sliding" | "dropping" | "suspend" | undefined
readonly replay?: number | undefinedreplay?: number | undefined
readonly idleTimeToLive?: Duration.Input | undefinedidleTimeToLive?: import DurationDuration.type Duration.Input = /*unresolved*/ anyInput | undefined
}
): import EffectEffect.type Effect.Effect = /*unresolved*/ anyEffect<interface Stream<out A, out E = never, out R = never>A Stream<A, E, R> describes a program that can emit many A values, fail
with E, and require R.
Details
Streams are pull-based with backpressure and emit chunks to amortize effect
evaluation. They support monadic composition and error handling similar to
Effect, adapted for multiple values.
Example (Creating and consuming streams)
import { Console, Effect, Stream } from "effect"
const program = Effect.gen(function*() {
yield* Stream.make(1, 2, 3).pipe(
Stream.map((n) => n * 2),
Stream.runForEach((n) => Console.log(n))
)
})
Effect.runPromise(program)
// Output:
// 2
// 4
// 6
Stream<function (type parameter) A in <A, E, R>(self: Stream<A, E, R>, options: {
readonly capacity: "unbounded";
readonly replay?: number | undefined;
readonly idleTimeToLive?: Duration.Input | undefined;
} | {
readonly capacity: number;
readonly strategy?: "sliding" | "dropping" | "suspend" | undefined;
readonly replay?: number | undefined;
readonly idleTimeToLive?: Duration.Input | undefined;
}): Effect.Effect<Stream<A, E>, never, Scope.Scope | R>
A, function (type parameter) E in <A, E, R>(self: Stream<A, E, R>, options: {
readonly capacity: "unbounded";
readonly replay?: number | undefined;
readonly idleTimeToLive?: Duration.Input | undefined;
} | {
readonly capacity: number;
readonly strategy?: "sliding" | "dropping" | "suspend" | undefined;
readonly replay?: number | undefined;
readonly idleTimeToLive?: Duration.Input | undefined;
}): Effect.Effect<Stream<A, E>, never, Scope.Scope | R>
E>, never, import ScopeScope.type Scope.Scope = /*unresolved*/ anyScope | function (type parameter) R in <A, E, R>(self: Stream<A, E, R>, options: {
readonly capacity: "unbounded";
readonly replay?: number | undefined;
readonly idleTimeToLive?: Duration.Input | undefined;
} | {
readonly capacity: number;
readonly strategy?: "sliding" | "dropping" | "suspend" | undefined;
readonly replay?: number | undefined;
readonly idleTimeToLive?: Duration.Input | undefined;
}): Effect.Effect<Stream<A, E>, never, Scope.Scope | R>
R>
} = import dualdual(2, <function (type parameter) A in <A, E, R>(self: Stream<A, E, R>, options: {
readonly capacity: "unbounded";
readonly replay?: number | undefined;
readonly idleTimeToLive?: Duration.Input | undefined;
} | {
readonly capacity: number;
readonly strategy?: "sliding" | "dropping" | "suspend" | undefined;
readonly replay?: number | undefined;
readonly idleTimeToLive?: Duration.Input | undefined;
}): Effect.Effect<Stream<A, E>, never, Scope.Scope | R>
A, function (type parameter) E in <A, E, R>(self: Stream<A, E, R>, options: {
readonly capacity: "unbounded";
readonly replay?: number | undefined;
readonly idleTimeToLive?: Duration.Input | undefined;
} | {
readonly capacity: number;
readonly strategy?: "sliding" | "dropping" | "suspend" | undefined;
readonly replay?: number | undefined;
readonly idleTimeToLive?: Duration.Input | undefined;
}): Effect.Effect<Stream<A, E>, never, Scope.Scope | R>
E, function (type parameter) R in <A, E, R>(self: Stream<A, E, R>, options: {
readonly capacity: "unbounded";
readonly replay?: number | undefined;
readonly idleTimeToLive?: Duration.Input | undefined;
} | {
readonly capacity: number;
readonly strategy?: "sliding" | "dropping" | "suspend" | undefined;
readonly replay?: number | undefined;
readonly idleTimeToLive?: Duration.Input | undefined;
}): Effect.Effect<Stream<A, E>, never, Scope.Scope | R>
R>(
self: Stream<A, E, R>(parameter) self: {
channel: Channel.Channel<Arr.NonEmptyReadonlyArray<A>, E, void, unknown, unknown, unknown, R>;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
self: interface Stream<out A, out E = never, out R = never>A Stream<A, E, R> describes a program that can emit many A values, fail
with E, and require R.
Details
Streams are pull-based with backpressure and emit chunks to amortize effect
evaluation. They support monadic composition and error handling similar to
Effect, adapted for multiple values.
Example (Creating and consuming streams)
import { Console, Effect, Stream } from "effect"
const program = Effect.gen(function*() {
yield* Stream.make(1, 2, 3).pipe(
Stream.map((n) => n * 2),
Stream.runForEach((n) => Console.log(n))
)
})
Effect.runPromise(program)
// Output:
// 2
// 4
// 6
Stream<function (type parameter) A in <A, E, R>(self: Stream<A, E, R>, options: {
readonly capacity: "unbounded";
readonly replay?: number | undefined;
readonly idleTimeToLive?: Duration.Input | undefined;
} | {
readonly capacity: number;
readonly strategy?: "sliding" | "dropping" | "suspend" | undefined;
readonly replay?: number | undefined;
readonly idleTimeToLive?: Duration.Input | undefined;
}): Effect.Effect<Stream<A, E>, never, Scope.Scope | R>
A, function (type parameter) E in <A, E, R>(self: Stream<A, E, R>, options: {
readonly capacity: "unbounded";
readonly replay?: number | undefined;
readonly idleTimeToLive?: Duration.Input | undefined;
} | {
readonly capacity: number;
readonly strategy?: "sliding" | "dropping" | "suspend" | undefined;
readonly replay?: number | undefined;
readonly idleTimeToLive?: Duration.Input | undefined;
}): Effect.Effect<Stream<A, E>, never, Scope.Scope | R>
E, function (type parameter) R in <A, E, R>(self: Stream<A, E, R>, options: {
readonly capacity: "unbounded";
readonly replay?: number | undefined;
readonly idleTimeToLive?: Duration.Input | undefined;
} | {
readonly capacity: number;
readonly strategy?: "sliding" | "dropping" | "suspend" | undefined;
readonly replay?: number | undefined;
readonly idleTimeToLive?: Duration.Input | undefined;
}): Effect.Effect<Stream<A, E>, never, Scope.Scope | R>
R>,
options: | {
readonly capacity: "unbounded"
readonly replay?: number | undefined
readonly idleTimeToLive?:
| Duration.Input
| undefined
}
| {
readonly capacity: number
readonly strategy?:
| "sliding"
| "dropping"
| "suspend"
| undefined
readonly replay?: number | undefined
readonly idleTimeToLive?:
| Duration.Input
| undefined
}
options: {
readonly capacity: "unbounded"capacity: "unbounded"
readonly replay?: number | undefinedreplay?: number | undefined
readonly idleTimeToLive?: Duration.Input | undefinedidleTimeToLive?: import DurationDuration.type Duration.Input = /*unresolved*/ anyInput | undefined
} | {
readonly capacity: numbercapacity: number
readonly strategy?: | "suspend"
| "sliding"
| "dropping"
| undefined
strategy?: "sliding" | "dropping" | "suspend" | undefined
readonly replay?: number | undefinedreplay?: number | undefined
readonly idleTimeToLive?: Duration.Input | undefinedidleTimeToLive?: import DurationDuration.type Duration.Input = /*unresolved*/ anyInput | undefined
}
): import EffectEffect.type Effect.Effect = /*unresolved*/ anyEffect<interface Stream<out A, out E = never, out R = never>A Stream<A, E, R> describes a program that can emit many A values, fail
with E, and require R.
Details
Streams are pull-based with backpressure and emit chunks to amortize effect
evaluation. They support monadic composition and error handling similar to
Effect, adapted for multiple values.
Example (Creating and consuming streams)
import { Console, Effect, Stream } from "effect"
const program = Effect.gen(function*() {
yield* Stream.make(1, 2, 3).pipe(
Stream.map((n) => n * 2),
Stream.runForEach((n) => Console.log(n))
)
})
Effect.runPromise(program)
// Output:
// 2
// 4
// 6
Stream<function (type parameter) A in <A, E, R>(self: Stream<A, E, R>, options: {
readonly capacity: "unbounded";
readonly replay?: number | undefined;
readonly idleTimeToLive?: Duration.Input | undefined;
} | {
readonly capacity: number;
readonly strategy?: "sliding" | "dropping" | "suspend" | undefined;
readonly replay?: number | undefined;
readonly idleTimeToLive?: Duration.Input | undefined;
}): Effect.Effect<Stream<A, E>, never, Scope.Scope | R>
A, function (type parameter) E in <A, E, R>(self: Stream<A, E, R>, options: {
readonly capacity: "unbounded";
readonly replay?: number | undefined;
readonly idleTimeToLive?: Duration.Input | undefined;
} | {
readonly capacity: number;
readonly strategy?: "sliding" | "dropping" | "suspend" | undefined;
readonly replay?: number | undefined;
readonly idleTimeToLive?: Duration.Input | undefined;
}): Effect.Effect<Stream<A, E>, never, Scope.Scope | R>
E>, never, import ScopeScope.type Scope.Scope = /*unresolved*/ anyScope | function (type parameter) R in <A, E, R>(self: Stream<A, E, R>, options: {
readonly capacity: "unbounded";
readonly replay?: number | undefined;
readonly idleTimeToLive?: Duration.Input | undefined;
} | {
readonly capacity: number;
readonly strategy?: "sliding" | "dropping" | "suspend" | undefined;
readonly replay?: number | undefined;
readonly idleTimeToLive?: Duration.Input | undefined;
}): Effect.Effect<Stream<A, E>, never, Scope.Scope | R>
R> =>
import EffectEffect.map(
import RcRefRcRef.make({
acquire: Effect.Effect<
Stream<A, E, never>,
never,
Scope.Scope | R
>
(property) acquire: {
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
toString: () => string;
toJSON: () => unknown;
}
acquire: const broadcast: {
(
options:
| {
readonly capacity: "unbounded"
readonly replay?: number | undefined
}
| {
readonly capacity: number
readonly strategy?:
| "sliding"
| "dropping"
| "suspend"
| undefined
readonly replay?: number | undefined
}
): <A, E, R>(
self: Stream<A, E, R>
) => Effect.Effect<
Stream<A, E>,
never,
Scope.Scope | R
>
<A, E, R>(
self: Stream<A, E, R>,
options:
| {
readonly capacity: "unbounded"
readonly replay?: number | undefined
}
| {
readonly capacity: number
readonly strategy?:
| "sliding"
| "dropping"
| "suspend"
| undefined
readonly replay?: number | undefined
}
): Effect.Effect<
Stream<A, E>,
never,
Scope.Scope | R
>
}
broadcast(self: Stream<A, E, R>(parameter) self: {
channel: Channel.Channel<Arr.NonEmptyReadonlyArray<A>, E, void, unknown, unknown, unknown, R>;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
self, options: | {
readonly capacity: "unbounded"
readonly replay?: number | undefined
readonly idleTimeToLive?:
| Duration.Input
| undefined
}
| {
readonly capacity: number
readonly strategy?:
| "sliding"
| "dropping"
| "suspend"
| undefined
readonly replay?: number | undefined
readonly idleTimeToLive?:
| Duration.Input
| undefined
}
options),
idleTimeToLive: anyidleTimeToLive: options: | {
readonly capacity: "unbounded"
readonly replay?: number | undefined
readonly idleTimeToLive?:
| Duration.Input
| undefined
}
| {
readonly capacity: number
readonly strategy?:
| "sliding"
| "dropping"
| "suspend"
| undefined
readonly replay?: number | undefined
readonly idleTimeToLive?:
| Duration.Input
| undefined
}
options.idleTimeToLive?: anyidleTimeToLive
}),
(ref: RcRef.RcRef<Stream<A, E, never>, never>(parameter) ref: {
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
ref) => const unwrap: <A, E2, R2, E, R>(
effect: Effect.Effect<Stream<A, E2, R2>, E, R>
) => Stream<
A,
E | E2,
R2 | Exclude<R, Scope.Scope>
>
Creates a stream produced from an Effect.
Example (Unwrapping a stream effect)
import { Console, Effect, Stream } from "effect"
const effect = Effect.succeed(Stream.make(1, 2, 3))
const stream = Stream.unwrap(effect)
const program = Effect.gen(function*() {
const chunk = yield* Stream.runCollect(stream)
yield* Console.log(chunk)
})
// [1, 2, 3]
unwrap(import RcRefRcRef.get(ref: RcRef.RcRef<Stream<A, E, never>, never>(parameter) ref: {
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
ref))
))