interface / @typed/fx

Push.Push

A bidirectional value that is both a Sink<A, E, R> and an Fx<B, E2, R2>.

Calling onSuccess or onFailure sends exactly one input notification to the wrapped Sink. The returned Effect is the acknowledgment: a producer that runs and awaits it waits for the consumer callback to finish. Push adds no queue, buffering, replay, or demand protocol of its own.

Running the Fx side preserves that Fx’s cardinality and ordering. Input values are not automatically forwarded to the output; any relationship between the two sides belongs to the supplied Sink and Fx (for example, a shared Subject).

Package version
2.0.0-beta.7
Category
Bidirectional contracts
Since
1.0.0
Member of
Push

Import

import { Push } from "@typed/fx";

Access Push.Push through the imported Push export. Its declaration below describes the member.

This public exposure is a re-export. Its import path is supported; declaration documentation is shared with the other public exposures below.

Signatures

export interface Push<in A, in E, out R, out B, out E2, out R2> extends Sink.Sink<A, E, R>, Fx.Fx<B, E2, R2> {
}
export declare namespace Push {
    type Any = Push<any, any, any, any, any, any>;
    interface Service<Self, Id extends string, A, E, B, E2> extends Push<A, E, Self, B, E2, Self> {
        readonly id: Id;
        readonly service: Context.Service<Self, Push<A, E, never, B, E2, never>>;
        readonly make: <R = never, R2 = never>(sink: Sink.Sink<A, E, R>, fx: Fx.Fx<B, E2, R2>) => Layer.Layer<Self, never, Exclude<R | R2, Scope.Scope>>;
    }
    interface Class<Self, Id extends string, A, E, B, E2> extends Service<Self, Id, A, E, B, E2> {
        new (): Service<Self, Id, A, E, B, E2>;
    }
}

Why

Interactive systems often need one value for sending work and another for observing results. Keeping both directions in one Effect-native value makes the boundary composable without pretending that commands and events are the same stream.

Ownership and lifetime

Constructing a Push starts no work and allocates no fiber or Scope. Each input callback runs in the caller’s fiber with the Sink’s R requirements. Running the output side follows the wrapped Fx’s Scope, interruption, finalizers, E2 failures, and R2 services; Push does not change them.

Examples

import { Effect } from "effect"
import * as Push from "@typed/fx/Push"
import { Fx } from "@typed/fx"
import * as Sink from "@typed/fx/Sink"

const program = Effect.gen(function* () {
  // Create a Push that accepts numbers and emits strings
  const push = Push.make(
    Sink.make(
      (cause) => Effect.sync(() => console.log("Error:", cause)),
      (value) => Effect.sync(() => console.log("Received:", value))
    ),
    Fx.succeed("Hello")
  )

  // Push a value to the sink
  yield* push.onSuccess(42)
  // Output: "Received: 42"

  // Observe the Fx output
  yield* Fx.observe(push, (value) =>
    Effect.sync(() => console.log("Emitted:", value))
  )
  // Output: "Emitted: Hello"
})

Members

  • Push.Push.Any

    Matches any Push when its six channel types are intentionally unknown.

  • Push.Push.Class

    Constructable static type produced by Push.Service.

  • Push.Push.Service

    The static and Effect service surface returned by Push.Service.

    Service lookup supplies the same bidirectional value to run, onSuccess, and onFailure. The Self service appears in both required-service channels; the installed Push itself has those requirements captured by its Layer.

Other public imports

These import paths expose the same declaration. Each page retains its own public name and signature.

Source