Hyperlinkv0.8.0-beta.28

Stream

Stream.takeWhileFilterconsteffect/Stream.ts:6589
<A, B, X>(f: Filter.Filter<NoInfer<A>, B, X>): <E, R>(
  self: Stream<A, E, R>
) => Stream<B, E, R>
<A, E, R, B, X>(
  self: Stream<A, E, R>,
  f: Filter.Filter<NoInfer<A>, B, X>
): Stream<B, E, R>

Takes the longest initial prefix accepted by a Filter and emits the filter's success values.

When to use

Use to keep the leading stream elements that a Filter accepts, emit the filter's success values, and stop at the first filter failure.

Details

The stream stops at the first Result.fail returned by the filter.

Source effect/Stream.ts:658930 lines
export const takeWhileFilter: {
  <A, B, X>(f: Filter.Filter<NoInfer<A>, B, X>): <E, R>(self: Stream<A, E, R>) => Stream<B, E, R>
  <A, E, R, B, X>(self: Stream<A, E, R>, f: Filter.Filter<NoInfer<A>, B, X>): Stream<B, E, R>
} = dual(
  2,
  <A, E, R, B, X>(
    self: Stream<A, E, R>,
    filter: Filter.Filter<NoInfer<A>, B, X>
  ): Stream<B, E, R> =>
    transformPull(self, (pull, _scope) =>
      Effect.sync(() => {
        let done = false
        const pump: Pull.Pull<Arr.NonEmptyReadonlyArray<B>, E, void, R> = Effect.flatMap(
          Effect.suspend(() => done ? Cause.done() : pull),
          (chunk) => {
            const out: Array<B> = []
            for (let j = 0; j < chunk.length; j++) {
              const result = filter(chunk[j])
              if (Result.isFailure(result)) {
                done = true
                break
              }
              out.push(result.success)
            }
            return Arr.isReadonlyArrayNonEmpty(out) ? Effect.succeed(out) : done ? Cause.done() : pump
          }
        )
        return pump
      }))
)