Hyperlinkv0.8.0-beta.28

Stream

Stream.tapSinkconsteffect/Stream.ts:2338
<A, E2, R2>(sink: 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.Sink<unknown, A, unknown, E2, R2>
): Stream<A, E | E2, R | R2>

Runs a sink for all stream elements while still emitting them downstream.

Example (Tapping values with a sink)

import { Console, Effect, Ref, Sink, Stream } from "effect"

const program = Effect.gen(function*() {
  const seen = yield* Ref.make<Array<number>>([])
  const sink = Sink.forEach((value: number) =>
    Ref.update(seen, (items) => [...items, value])
  )
  const result = yield* Stream.make(1, 2, 3).pipe(
    Stream.tapSink(sink),
    Stream.runCollect
  )
  const tapped = yield* Ref.get(seen)
  yield* Console.log(tapped)
  yield* Console.log(result)
})

Effect.runPromise(program)
// Output: [1, 2, 3]
// Output: [1, 2, 3]
sequencing
Source effect/Stream.ts:233867 lines
export const tapSink: {
  <A, E2, R2>(sink: 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.Sink<unknown, A, unknown, E2, R2>): Stream<A, E | E2, R | R2>
} = dual(
  2,
  <A, E, R, E2, R2>(
    self: Stream<A, E, R>,
    sink: Sink.Sink<unknown, A, unknown, E2, R2>
  ): Stream<A, E | E2, R | R2> =>
    transformPullBracket(
      self,
      Effect.fnUntraced(function*(pull, _, scope) {
        const upstreamLatch = Latch.makeUnsafe()
        const sinkLatch = Latch.makeUnsafe()
        let chunk: Arr.NonEmptyReadonlyArray<A> | undefined = undefined
        let causeSink: Cause.Cause<E2> | undefined = undefined
        let sinkDone = false
        let streamDone = false

        const sinkUpstream = upstreamLatch.whenOpen(Effect.suspend(() => {
          if (chunk) {
            const arr = chunk!
            chunk = undefined
            if (!streamDone) upstreamLatch.closeUnsafe()
            return Effect.as(sinkLatch.open, arr)
          }
          return Cause.done()
        }))

        yield* Effect.suspend(() => sink.transform(sinkUpstream, scope)).pipe(
          (eff) =>
            Effect.onExitPrimitive(eff, (exit) => {
              sinkDone = true
              if (Exit.isFailure(exit)) {
                causeSink = exit.cause
              }
              return sinkLatch.open
            }, true),
          Effect.forkIn(scope)
        )

        const pullAndOffer = pull.pipe(
          Effect.flatMap((chunk_) => {
            chunk = chunk_
            sinkLatch.closeUnsafe()
            upstreamLatch.openUnsafe()
            return Effect.as(sinkLatch.await, chunk_)
          }),
          Pull.catchDone(() => {
            streamDone = true
            sinkLatch.closeUnsafe()
            upstreamLatch.openUnsafe()
            return Effect.flatMap(sinkLatch.await, () => Cause.done())
          })
        )

        return Effect.suspend((): Pull.Pull<Arr.NonEmptyReadonlyArray<A>, E | E2, void, R> => {
          if (causeSink) {
            return Effect.failCause(causeSink)
          } else if (sinkDone) {
            return pull
          }
          return pullAndOffer
        })
      })
    )
)