MCPcopy Create free account
hub / github.com/effect-app/libs / SemaphoreImpl

Class SemaphoreImpl

repos/effect/packages/effect/src/Semaphore.ts:192–294  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

190 * // Only 3 tasks can run concurrently
191 * const program = Effect.all([
192 * task(1),
193 * task(2),
194 * task(3),
195 * task(4),
196 * task(5)
197 * ], { concurrency: "unbounded" })
198 *
199 * await Effect.runPromise(program) // => [1, 2, 3, 4, 5]
200 * ```
201 *
202 * @category constructors
203 * @since 4.0.0
204 */
205export const makeUnsafe = (permits: number): Semaphore => new SemaphoreImpl(permits)
206
207const waitForPermits = <A, E, R>(
208 self: SemaphoreImpl,
209 n: number,
210 effect: Effect.Effect<A, E, R>
211): Effect.Effect<A, E, R> =>
212 internal.callback((resume) => {
213 if (self.free >= n) return resume(effect)
214 const observer = () => {
215 if (self.free < n) return
216 self.waiters.delete(observer)
217 resume(effect)
218 }
219 self.waiters.add(observer)
220 return internal.sync(() => {
221 self.waiters.delete(observer)
222 })
223 })
224
225class SemaphoreImpl implements Semaphore {
226 public waiters = new Set<() => void>()
227 public taken = 0
228 public permits: number
229
230 constructor(permits: number) {
231 this.permits = permits
232 }
233
234 get free() {
235 return this.permits - this.taken
236 }
237
238 take(n: number): Effect.Effect<number> {
239 const take: Effect.Effect<number> = internal.suspend(() => {
240 if (this.free < n) {
241 return waitForPermits(this, n, take)
242 }
243 this.taken += n
244 return internal.succeed(n)
245 })
246 return take
247 }
248
249 takeIfAvailable(n: number): Effect.Effect<boolean> {

Callers

nothing calls this directly

Calls 1

withPermitsMethod · 0.95

Tested by

no test coverage detected

Used in the wild real call sites across dependent graphs

searching dependent graphs…