Hyperlinkv0.8.0-beta.28

Channel

Channel.embedInputconsteffect/Channel.ts:6692
<InElem, InErr, InDone, R>(
  input: (
    upstream: Pull.Pull<InElem, InErr, InDone>
  ) => Effect.Effect<void, never, R>
): <OutElem, OutErr, OutDone, Env>(
  self: Channel<OutElem, OutErr, OutDone, unknown, unknown, unknown, Env>
) => Channel<OutElem, OutErr, OutDone, InElem, InErr, InDone, Env | R>
<OutElem, OutErr, OutDone, Env, InErr, InElem, InDone, R>(
  self: Channel<OutElem, OutErr, OutDone, unknown, unknown, unknown, Env>,
  input: (
    upstream: Pull.Pull<InElem, InErr, InDone>
  ) => Effect.Effect<void, never, R>
): Channel<OutElem, OutErr, OutDone, InElem, InErr, InDone, Env | R>

Runs an input handler against the upstream pull while the wrapped channel runs without receiving upstream input directly.

Details

The input handler is forked in the channel scope. The wrapped channel is run with an already-completed input.

Example (Embedding custom input handling)

import { Channel, Effect } from "effect"

// Create a base channel
const baseChannel = Channel.fromIterable([1, 2, 3])

// Drain the embedded input while the base channel runs
const embeddedChannel = Channel.embedInput(
  baseChannel,
  (upstream) =>
    upstream.pipe(
      Effect.tap((message) =>
        Effect.sync(() => console.log(message))
      ),
      Effect.forever,
      Effect.ignore
    )
)
sequencing
Source effect/Channel.ts:669229 lines
export const embedInput: {
  <InElem, InErr, InDone, R>(
    input: (
      upstream: Pull.Pull<InElem, InErr, InDone>
    ) => Effect.Effect<void, never, R>
  ): <OutElem, OutErr, OutDone, Env>(
    self: Channel<OutElem, OutErr, OutDone, unknown, unknown, unknown, Env>
  ) => Channel<OutElem, OutErr, OutDone, InElem, InErr, InDone, Env | R>
  <OutElem, OutErr, OutDone, Env, InErr, InElem, InDone, R>(
    self: Channel<OutElem, OutErr, OutDone, unknown, unknown, unknown, Env>,
    input: (
      upstream: Pull.Pull<InElem, InErr, InDone>
    ) => Effect.Effect<void, never, R>
  ): Channel<OutElem, OutErr, OutDone, InElem, InErr, InDone, Env | R>
} = dual(
  2,
  <OutElem, OutErr, OutDone, Env, InErr, InElem, InDone, R>(
    self: Channel<OutElem, OutErr, OutDone, unknown, unknown, unknown, Env>,
    input: (
      upstream: Pull.Pull<InElem, InErr, InDone>
    ) => Effect.Effect<void, never, R>
  ): Channel<OutElem, OutErr, OutDone, InElem, InErr, InDone, Env | R> =>
    fromTransformBracket((upstream, scope, forkedScope) =>
      Effect.andThen(
        Effect.forkIn(input(upstream), forkedScope),
        toTransform(self)(Cause.done(), scope)
      )
    )
)