| 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 | */ |
| 205 | export const makeUnsafe = (permits: number): Semaphore => new SemaphoreImpl(permits) |
| 206 | |
| 207 | const 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 | |
| 225 | class 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> { |
nothing calls this directly
no test coverage detected
searching dependent graphs…