Hyperlinkv0.8.0-beta.28

Stream

Stream.retryconsteffect/Stream.ts:6145
<E, X, E2, R2>(
  policy:
    | Schedule.Schedule<X, NoInfer<E>, E2, R2>
    | ((
        $: <SO, SE, SR>(
          _: Schedule.Schedule<SO, NoInfer<E>, SE, SR>
        ) => Schedule.Schedule<SO, E, SE, SR>
      ) => Schedule.Schedule<X, NoInfer<E>, E2, R2>)
): <A, R>(self: Stream<A, E, R>) => Stream<A, E | E2, R2 | R>
<A, E, R, X, E2, R2>(
  self: Stream<A, E, R>,
  policy:
    | Schedule.Schedule<X, NoInfer<E>, E2, R2>
    | ((
        $: <SO, SE, SR>(
          _: Schedule.Schedule<SO, NoInfer<E>, SE, SR>
        ) => Schedule.Schedule<SO, E, SE, SR>
      ) => Schedule.Schedule<X, NoInfer<E>, E2, R2>)
): Stream<A, E | E2, R2 | R>

Retries the stream according to the given schedule when it fails.

Details

This retries the entire stream, so will re-execute all of the stream's acquire operations.

The schedule is reset as soon as the first element passes through the stream again.

Example (Retrying stream failures)

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

const program = Effect.gen(function*() {
  const values = yield* Stream.make(1).pipe(
    Stream.concat(Stream.fail("boom")),
    Stream.retry(Schedule.recurs(1)),
    Stream.take(2),
    Stream.runCollect
  )

  yield* Console.log(values)
})

Effect.runPromise(program)
// Output: [ 1, 1 ]
error handling
Source effect/Stream.ts:614527 lines
export const retry: {
  <E, X, E2, R2>(
    policy:
      | Schedule.Schedule<X, NoInfer<E>, E2, R2>
      | ((
        $: <SO, SE, SR>(_: Schedule.Schedule<SO, NoInfer<E>, SE, SR>) => Schedule.Schedule<SO, E, SE, SR>
      ) => Schedule.Schedule<X, NoInfer<E>, E2, R2>)
  ): <A, R>(self: Stream<A, E, R>) => Stream<A, E | E2, R2 | R>
  <A, E, R, X, E2, R2>(
    self: Stream<A, E, R>,
    policy:
      | Schedule.Schedule<X, NoInfer<E>, E2, R2>
      | ((
        $: <SO, SE, SR>(_: Schedule.Schedule<SO, NoInfer<E>, SE, SR>) => Schedule.Schedule<SO, E, SE, SR>
      ) => Schedule.Schedule<X, NoInfer<E>, E2, R2>)
  ): Stream<A, E | E2, R2 | R>
} = dual(
  2,
  <A, E, R, X, E2, R2>(
    self: Stream<A, E, R>,
    policy:
      | Schedule.Schedule<X, NoInfer<E>, E2, R2>
      | ((
        $: <SO, SE, SR>(_: Schedule.Schedule<SO, NoInfer<E>, SE, SR>) => Schedule.Schedule<SO, E, SE, SR>
      ) => Schedule.Schedule<X, NoInfer<E>, E2, R2>)
  ): Stream<A, E | E2, R2 | R> => fromChannel(Channel.retry(self.channel, policy))
)
Referenced by 1 symbols