(self: Subscription<A>)
| 1210 | }) |
| 1211 | |
| 1212 | const pollForItem = <A>(self: Subscription<A>) => { |
| 1213 | const deferred = Deferred.makeUnsafe<A>() |
| 1214 | let set = self.subscribers.get(self.subscription) |
| 1215 | if (!set) { |
| 1216 | set = new Set() |
| 1217 | self.subscribers.set(self.subscription, set) |
| 1218 | } |
| 1219 | set.add(self.pollers) |
| 1220 | MutableList.append(self.pollers, deferred) |
| 1221 | self.strategy.completePollersUnsafe( |
| 1222 | self.pubsub, |
| 1223 | self.subscribers, |
| 1224 | self.subscription, |
| 1225 | self.pollers |
| 1226 | ) |
| 1227 | return Effect.onInterrupt( |
| 1228 | Deferred.await(deferred), |
| 1229 | () => { |
| 1230 | MutableList.remove(self.pollers, deferred) |
| 1231 | return Effect.void |
| 1232 | } |
| 1233 | ) |
| 1234 | } |
| 1235 | |
| 1236 | /** |
| 1237 | * Takes up to the specified number of messages from the subscription without suspending. |
no test coverage detected