(
register: (emit: {
readonly single: (item: A) => void
readonly array: (arr: ReadonlyArray<A>) => void
readonly fail: (error: E) => void
readonly end: () => void
}) => Effect.Effect<
{
onPause(): void
onResume(): void
},
E,
R | Scope.Scope
>,
bufferSize = 128
)
| 25 | * @since 4.0.0 |
| 26 | */ |
| 27 | export const asyncPauseResume = <A, E = never, R = never>( |
| 28 | register: (emit: { |
| 29 | readonly single: (item: A) => void |
| 30 | readonly array: (arr: ReadonlyArray<A>) => void |
| 31 | readonly fail: (error: E) => void |
| 32 | readonly end: () => void |
| 33 | }) => Effect.Effect< |
| 34 | { |
| 35 | onPause(): void |
| 36 | onResume(): void |
| 37 | }, |
| 38 | E, |
| 39 | R | Scope.Scope |
| 40 | >, |
| 41 | bufferSize = 128 |
| 42 | ): Stream.Stream<A, E, R> => |
| 43 | Stream.callback<A, E, R>((queue) => |
| 44 | Effect.suspend(() => { |
| 45 | let cbs!: { |
| 46 | onPause(): void |
| 47 | onResume(): void |
| 48 | } |
| 49 | |
| 50 | let paused = false |
| 51 | const offer = (arr: ReadonlyArray<A>) => { |
| 52 | if (arr.length === 0) return |
| 53 | const isFull = Queue.isFullUnsafe(queue) |
| 54 | if (!isFull || (isFull && paused)) { |
| 55 | return Effect.runFork(Queue.offerAll(queue, arr)) |
| 56 | } |
| 57 | paused = true |
| 58 | cbs.onPause() |
| 59 | return Queue.offerAll(queue, arr).pipe( |
| 60 | Effect.tap(() => |
| 61 | Effect.sync(() => { |
| 62 | cbs.onResume() |
| 63 | paused = false |
| 64 | }) |
| 65 | ), |
| 66 | Effect.runFork |
| 67 | ) |
| 68 | } |
| 69 | |
| 70 | return Effect.map( |
| 71 | register({ |
| 72 | single: (item) => offer([item]), |
| 73 | array: (chunk) => offer(chunk), |
| 74 | fail: (error) => Queue.failCauseUnsafe(queue as any, Cause.fail(error)), |
| 75 | end: () => Queue.endUnsafe(queue as any) |
| 76 | }), |
| 77 | (_) => { |
| 78 | cbs = _ |
| 79 | } |
| 80 | ) |
| 81 | }), { bufferSize }) |
no test coverage detected