Hyperlinkv0.8.0-beta.28

Stream

Stream.zipLatestAllconsteffect/Stream.ts:4034
<T extends ReadonlyArray<Stream<any, any, any>>>(...streams: T): Stream<
  [T[number]] extends [never]
    ? never
    : {
        [K in keyof T]: T[K] extends Stream<infer A, infer _E, infer _R>
          ? A
          : never
      },
  [T[number]] extends [never]
    ? never
    : T[number] extends Stream<infer _A, infer _E, infer _R>
    ? _E
    : never,
  [T[number]] extends [never]
    ? never
    : T[number] extends Stream<infer _A, infer _E, infer _R>
    ? _R
    : never
>

Zips multiple streams so that when a value is emitted by any stream, it is combined with the latest values from the other streams to produce a result.

When to use

Use when each stream should contribute its latest value after all streams have emitted at least once.

Gotchas

Note: tracking the latest value is done on a per-array basis. That means that emitted elements that are not the last value in arrays will never be used for zipping.

Example (Zipping latest values from many streams)

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

const stream = Stream.zipLatestAll(
  Stream.make(1, 2, 3).pipe(Stream.rechunk(1)),
  Stream.make("a", "b", "c").pipe(Stream.rechunk(1)),
  Stream.make(true, false, true).pipe(Stream.rechunk(1))
)

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

Effect.runPromise(program)
// Output: [ [ 1, "a", true ], [ 2, "a", true ], [ 3, "a", true ], [ 3, "b", true ], [ 3, "c", true ], [ 3, "c", false ], [ 3, "c", true ] ]
zipping
Source effect/Stream.ts:403438 lines
export const zipLatestAll = <T extends ReadonlyArray<Stream<any, any, any>>>(
  ...streams: T
): Stream<
  [T[number]] extends [never] ? never
    : { [K in keyof T]: T[K] extends Stream<infer A, infer _E, infer _R> ? A : never },
  [T[number]] extends [never] ? never : T[number] extends Stream<infer _A, infer _E, infer _R> ? _E : never,
  [T[number]] extends [never] ? never : T[number] extends Stream<infer _A, infer _E, infer _R> ? _R : never
> =>
  fromChannel(Channel.suspend(() => {
    const latest: Array<any> = []
    const emitted = new Set<number>()
    const readyLatch = Latch.makeUnsafe()
    return Channel.mergeAll(
      Channel.fromArray(
        streams.map((s, i) =>
          s.channel.pipe(
            Channel.flattenArray,
            Channel.mapEffect((a) => {
              latest[i] = a
              if (!emitted.has(i)) {
                emitted.add(i)
                if (emitted.size < streams.length) {
                  return readyLatch.await as Effect.Effect<undefined>
                }
                return Effect.as(readyLatch.open, Arr.of(latest.slice()))
              }
              return Effect.succeed(Arr.of(latest.slice()))
            }),
            Channel.filter(isNotUndefined)
          )
        )
      ),
      {
        concurrency: "unbounded",
        bufferSize: 0
      }
    )
  })) as any
Referenced by 2 symbols