Skip to content

Commit

Permalink
feat(Stream): zipLatestAll (#2832)
Browse files Browse the repository at this point in the history
Co-authored-by: Tim <hello@timsmart.co>
  • Loading branch information
dilame and tim-smart committed Jun 6, 2024
1 parent 64565db commit 1f4ac00
Show file tree
Hide file tree
Showing 5 changed files with 100 additions and 0 deletions.
5 changes: 5 additions & 0 deletions .changeset/calm-bags-admire.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
"effect": minor
---

add `Stream.zipLatestAll` api
13 changes: 13 additions & 0 deletions packages/effect/dtslint/Stream.ts
Original file line number Diff line number Diff line change
Expand Up @@ -236,3 +236,16 @@ pipe(
_scope // $ExpectType { a: number; b: string; }
) => true)
)

// -------------------------------------------------------------------------------------
// zipLatestAll
// -------------------------------------------------------------------------------------

// $ExpectType Stream<never, never, never>
Stream.zipLatestAll()

// $ExpectType Stream<[number, string | number], never, never>
Stream.zipLatestAll(numbers, numbersOrStrings)

// $ExpectType Stream<[number, string | number, never], Error, never>
Stream.zipLatestAll(numbers, numbersOrStrings, Stream.fail(new Error("")))
39 changes: 39 additions & 0 deletions packages/effect/src/Stream.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4350,6 +4350,45 @@ export const zipLatest: {
<A, E, R, A2, E2, R2>(self: Stream<A, E, R>, that: Stream<A2, E2, R2>): Stream<[A, A2], E | E2, R | R2>
} = internal.zipLatest

/**
* Zips multiple streams so that when a value is emitted by any of the streams,
* it is combined with the latest values from the other streams to produce a result.
*
* Note: tracking the latest value is done on a per-chunk basis. That means
* that emitted elements that are not the last value in chunks will never be
* used for zipping.
*
* @example
* import { Stream, Schedule, Console, Effect } from "effect"
*
* const stream = Stream.zipLatestAll(
* Stream.fromSchedule(Schedule.spaced('1 millis')),
* Stream.fromSchedule(Schedule.spaced('2 millis')),
* Stream.fromSchedule(Schedule.spaced('4 millis')),
* ).pipe(Stream.take(6), Stream.tap(Console.log))
*
* Effect.runPromise(Stream.runDrain(stream))
* // Output:
* // [ 0, 0, 0 ]
* // [ 1, 0, 0 ]
* // [ 1, 1, 0 ]
* // [ 2, 1, 0 ]
* // [ 3, 1, 0 ]
* // [ 3, 1, 1 ]
* // .....
*
* @since 3.3.0
* @category zipping
*/
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
> = internal.zipLatestAll

/**
* Zips the two streams so that when a value is emitted by either of the two
* streams, it is combined with the latest value from the other stream to
Expand Down
21 changes: 21 additions & 0 deletions packages/effect/src/internal/stream.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7609,6 +7609,27 @@ export const zipLatest = dual<
): Stream.Stream<[A, A2], E2 | E, R2 | R> => pipe(self, zipLatestWith(that, (a, a2) => [a, a2]))
)

export const zipLatestAll = <T extends ReadonlyArray<Stream.Stream<any, any, any>>>(
...streams: T
): Stream.Stream<
[T[number]] extends [never] ? never
: { [K in keyof T]: T[K] extends Stream.Stream<infer A, infer _E, infer _R> ? A : never },
[T[number]] extends [never] ? never : T[number] extends Stream.Stream<infer _A, infer _E, infer _R> ? _E : never,
[T[number]] extends [never] ? never : T[number] extends Stream.Stream<infer _A, infer _E, infer _R> ? _R : never
> => {
if (streams.length === 0) {
return empty
} else if (streams.length === 1) {
return map(streams[0]!, (x) => [x]) as any
}
const [head, ...tail] = streams
return zipLatestWith(
head,
zipLatestAll(...tail),
(first, second) => [first, ...second]
) as any
}

/** @internal */
export const zipLatestWith = dual<
<A2, E2, R2, A, A3>(
Expand Down
22 changes: 22 additions & 0 deletions packages/effect/test/Stream/zipping.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -466,4 +466,26 @@ describe("Stream", () => {
)
assert.deepStrictEqual(Array.from(result1), Array.from(result2))
})))

it.effect("zipLatestAll", () =>
Effect.gen(function*($) {
const result = yield* $(
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))
),
Stream.runCollect,
Effect.map(Chunk.toReadonlyArray)
)
assert.deepStrictEqual(result, [
[1, "a", true],
[2, "a", true],
[3, "a", true],
[3, "b", true],
[3, "c", true],
[3, "c", false],
[3, "c", true]
])
}))
})

0 comments on commit 1f4ac00

Please sign in to comment.