Hyperlinkv0.8.0-beta.28

Stream

Stream.takeWhileconsteffect/Stream.ts:6536
<A, B extends A>(refinement: (a: NoInfer<A>, n: number) => a is B): <
  E,
  R
>(
  self: Stream<A, E, R>
) => Stream<B, E, R>
<A>(predicate: (a: NoInfer<A>, n: number) => boolean): <E, R>(
  self: Stream<A, E, R>
) => Stream<A, E, R>
<A, E, R, B extends A>(
  self: Stream<A, E, R>,
  refinement: (a: NoInfer<A>, n: number) => a is B
): Stream<B, E, R>
<A, E, R>(
  self: Stream<A, E, R>,
  predicate: (a: NoInfer<A>, n: number) => boolean
): Stream<A, E, R>

Takes the longest initial prefix of elements that satisfy the predicate.

Example (Taking while a predicate holds)

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

const stream = Stream.range(1, 5).pipe(
  Stream.takeWhile((n) => n % 3 !== 0)
)

const program = Effect.gen(function*() {
  const result = yield* Stream.runCollect(stream)
  yield* Console.log(result)
})

Effect.runPromise(program)
// Output: [ 1, 2 ]
filtering
Source effect/Stream.ts:653632 lines
export const takeWhile: {
  <A, B extends A>(refinement: (a: NoInfer<A>, n: number) => a is B): <E, R>(self: Stream<A, E, R>) => Stream<B, E, R>
  <A>(predicate: (a: NoInfer<A>, n: number) => boolean): <E, R>(self: Stream<A, E, R>) => Stream<A, E, R>
  <A, E, R, B extends A>(self: Stream<A, E, R>, refinement: (a: NoInfer<A>, n: number) => a is B): Stream<B, E, R>
  <A, E, R>(self: Stream<A, E, R>, predicate: (a: NoInfer<A>, n: number) => boolean): Stream<A, E, R>
} = dual(
  2,
  <A, E, R>(
    self: Stream<A, E, R>,
    predicate: (a: A, n: number) => boolean
  ): Stream<A, E, R> =>
    transformPull(self, (pull, _scope) =>
      Effect.sync(() => {
        let i = 0
        let done = false
        const pump: Pull.Pull<Arr.NonEmptyReadonlyArray<A>, E, void, R> = Effect.flatMap(
          Effect.suspend(() => done ? Cause.done() : pull),
          (chunk) => {
            const out: Array<A> = []
            for (let j = 0; j < chunk.length; j++) {
              if (!predicate(chunk[j], i++)) {
                done = true
                break
              }
              out.push(chunk[j])
            }
            return Arr.isReadonlyArrayNonEmpty(out) ? Effect.succeed(out) : done ? Cause.done() : pump
          }
        )
        return pump
      }))
)