MCPcopy Create free account
hub / github.com/Effect-TS/effect / asyncPauseResume

Function asyncPauseResume

packages/effect/src/unstable/sql/SqlStream.ts:27–81  ·  view source on GitHub ↗
(
  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
)

Source from the content-addressed store, hash-verified

25 * @since 4.0.0
26 */
27export 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 })

Callers 1

queryStreamFunction · 0.90

Calls 4

registerFunction · 0.85
offerFunction · 0.70
mapMethod · 0.45
failMethod · 0.45

Tested by

no test coverage detected