Stream.Stream<
DirectoryUpserted | DirectoryRemoved,
never,
Directory
>Live directory membership push — sugar over Directory.changes.
Subscribe from a node process to notice A→B dial swaps (DirectoryUpserted with
dialChanged: true) and rebind peer clients without restart.
yield* Lookup.changes.pipe(
Stream.filter((e) => e._tag === "DirectoryUpserted" && e.dialChanged),
Stream.runForEach((e) => Effect.logInfo(`dial moved ${e.entry.nodeKey}`)),
)export const const changes: Stream.Stream<
DirectoryChange,
never,
Directory
>
const changes: {
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; <…;
}
Live directory membership push — sugar over
Directory
.changes.
Subscribe from a node process to notice A→B dial swaps (DirectoryUpserted with
dialChanged: true) and rebind peer clients without restart.
yield* Lookup.changes.pipe(
Stream.filter((e) => e._tag === "DirectoryUpserted" && e.dialChanged),
Stream.runForEach((e) => Effect.logInfo(`dial moved ${e.entry.nodeKey}`)),
)
changes: import StreamStream.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<type DirectoryChange =
| DirectoryUpserted
| DirectoryRemoved
Membership-push event from
Directory.changes
.
Membership-push event type.
DirectoryChange, never, class Directoryclass Directory {
key: Identifier;
Service: {
advertise: (payload: { kind: 'Http' | 'WebSocket' | 'IpcSocket'; serves: ReadonlyArray<string>; nodeKey: string; url?: string | undefined; path?: string | undefined; onConflict?: 'livenessReplace' | 'askIncumbent' | 'reject' | 'inherit' | undefined }…;
unregister: (payload: { nodeKey: string; kind?: 'Http' | 'WebSocket' | 'IpcSocket' | undefined; url?: string | undefined; path?: string | undefined }) => Effect.Effect<boolean, never, never>;
nodesServing: (payload: { serviceKey: string }) => Effect.Effect<ReadonlyArray<DirectoryEntry>, never, never>;
changes: Stream.Stream<DirectoryUpserted | DirectoryRemoved, never, never>;
};
}
Lookup node directory — advertise / unregister / list by served HyperService key /
directorySpec.changes
membership push.
Directory> =
import StreamStream.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 EffectEffect.const map: {
<A, B>(f: (a: A) => B): <E, R>(
self: Effect<A, E, R>
) => Effect<B, E, R>
<A, E, R, B>(
self: Effect<A, E, R>,
f: (a: A) => B
): Effect<B, E, R>
}
Transforms the value inside an effect by applying a function to it.
When to use
Use to transform an effect's success value with a function that returns a
plain value, producing a new effect without changing the original effect's
typed error or context requirements.
Details
map takes a function and applies it to the value contained within an
effect, creating a new effect with the transformed value.
It's important to note that effects are immutable, meaning that the original
effect is not modified. Instead, a new effect is returned with the updated
value.
Example (Choosing map syntax variants)
import { Effect, pipe } from "effect"
const myEffect = Effect.succeed(1)
const transformation = (n: number) => n + 1
const mappedWithPipe = pipe(myEffect, Effect.map(transformation))
const mappedWithDataFirst = Effect.map(myEffect, transformation)
const mappedWithMethod = myEffect.pipe(Effect.map(transformation))
Example (Adding a service charge)
import { Effect, pipe } from "effect"
const addServiceCharge = (amount: number) => amount + 1
const fetchTransactionAmount = Effect.promise(() => Promise.resolve(100))
const finalAmount = pipe(
fetchTransactionAmount,
Effect.map(addServiceCharge)
)
Effect.runPromise(finalAmount).then(console.log)
// Output: 101
map(class Directoryclass Directory {
key: Identifier;
Service: {
advertise: (payload: { kind: 'Http' | 'WebSocket' | 'IpcSocket'; serves: ReadonlyArray<string>; nodeKey: string; url?: string | undefined; path?: string | undefined; onConflict?: 'livenessReplace' | 'askIncumbent' | 'reject' | 'inherit' | undefined }…;
unregister: (payload: { nodeKey: string; kind?: 'Http' | 'WebSocket' | 'IpcSocket' | undefined; url?: string | undefined; path?: string | undefined }) => Effect.Effect<boolean, never, never>;
nodesServing: (payload: { serviceKey: string }) => Effect.Effect<ReadonlyArray<DirectoryEntry>, never, never>;
changes: Stream.Stream<DirectoryUpserted | DirectoryRemoved, never, never>;
};
description: string | undefined;
of: (this: void, self: { readonly advertise: (payload: { kind: 'Http' | 'WebSocket' | 'IpcSocket'; serves: ReadonlyArray<string>; nodeKey: string; url?: string | undefined; path?: string | undefined; onConflict?: 'livenessReplace' | 'askIncumb…;
context: (self: { readonly advertise: (payload: { kind: 'Http' | 'WebSocket' | 'IpcSocket'; serves: ReadonlyArray<string>; nodeKey: string; url?: string | undefined; path?: string | undefined; onConflict?: 'livenessReplace' | 'askIncumbent' | 'reje…;
use: (f: (service: { readonly advertise: (payload: { kind: 'Http' | 'WebSocket' | 'IpcSocket'; serves: ReadonlyArray<string>; nodeKey: string; url?: string | undefined; path?: string | undefined; onConflict?: 'livenessReplace' | 'askIncumbent' …;
useSync: (f: (service: { readonly advertise: (payload: { kind: 'Http' | 'WebSocket' | 'IpcSocket'; serves: ReadonlyArray<string>; nodeKey: string; url?: string | undefined; path?: string | undefined; onConflict?: 'livenessReplace' | 'askIncumbent' …;
Identifier: Identifier;
stack: string | undefined;
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;
}
Lookup node directory — advertise / unregister / list by served HyperService key /
directorySpec.changes
membership push.
Directory, (dir: {
readonly advertise: (payload: {
kind: "Http" | "WebSocket" | "IpcSocket";
serves: readonly string[];
nodeKey: string;
url?: string | undefined;
path?: string | undefined;
onConflict?: "livenessReplace" | "askIncumbent" | "reject" | "inherit" | undefined;
}) => Effect.Effect<DirectoryEntry, IncumbentAlive, never>;
readonly unregister: (payload: {
nodeKey: string;
kind?: "Http" | "WebSocket" | "IpcSocket" | undefined;
url?: string | undefined;
path?: string | undefined;
}) => Effect.Effect<boolean, never, never>;
readonly nodesServing: (payload: {
...;
}) => Effect.Effect<...>;
readonly changes: Stream.Stream<...>;
}
dir) => dir: {
readonly advertise: (payload: {
kind: "Http" | "WebSocket" | "IpcSocket";
serves: readonly string[];
nodeKey: string;
url?: string | undefined;
path?: string | undefined;
onConflict?: "livenessReplace" | "askIncumbent" | "reject" | "inherit" | undefined;
}) => Effect.Effect<DirectoryEntry, IncumbentAlive, never>;
readonly unregister: (payload: {
nodeKey: string;
kind?: "Http" | "WebSocket" | "IpcSocket" | undefined;
url?: string | undefined;
path?: string | undefined;
}) => Effect.Effect<boolean, never, never>;
readonly nodesServing: (payload: {
...;
}) => Effect.Effect<...>;
readonly changes: Stream.Stream<...>;
}
dir.changes: Stream.Stream<
DirectoryUpserted | DirectoryRemoved,
never,
never
>
(property) changes: {
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; <…;
}
changes));