(
pubsub: PubSub.Atomic<A>,
subscribers: PubSub.Subscribers<A>,
elements: Iterable<A>,
isShutdown: MutableRef.MutableRef<boolean>
)
| 2370 | } |
| 2371 | |
| 2372 | handleSurplus( |
| 2373 | pubsub: PubSub.Atomic<A>, |
| 2374 | subscribers: PubSub.Subscribers<A>, |
| 2375 | elements: Iterable<A>, |
| 2376 | isShutdown: MutableRef.MutableRef<boolean> |
| 2377 | ): Effect.Effect<boolean> { |
| 2378 | return Effect.suspend(() => { |
| 2379 | const deferred = Deferred.makeUnsafe<boolean>() |
| 2380 | this.offerUnsafe(elements, deferred) |
| 2381 | this.onPubSubEmptySpaceUnsafe(pubsub, subscribers) |
| 2382 | this.completeSubscribersUnsafe(pubsub, subscribers) |
| 2383 | return (MutableRef.get(isShutdown) ? Effect.interrupt : Deferred.await(deferred)).pipe( |
| 2384 | Effect.onInterrupt(() => { |
| 2385 | this.removeUnsafe(deferred) |
| 2386 | return Effect.void |
| 2387 | }) |
| 2388 | ) |
| 2389 | }) |
| 2390 | } |
| 2391 | |
| 2392 | onPubSubEmptySpaceUnsafe( |
| 2393 | pubsub: PubSub.Atomic<A>, |
nothing calls this directly
no test coverage detected