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

Function getPoolItem

repos/effect/packages/effect/src/Pool.ts:429–465  ·  view source on GitHub ↗
(self: Pool<A, E>)

Source from the content-addressed store, hash-verified

427 }
428 return Effect.flatMap(getPoolItem(self), (item) => item.exit)
429 })
430
431const getPoolItem = <A, E>(self: Pool<A, E>): Effect.Effect<PoolItem<A, E>, never, Scope.Scope> =>
432 Effect.uninterruptibleMask((restore) =>
433 restore(self.state.semaphore.take(1)).pipe(
434 Effect.flatMap(() => Effect.scope),
435 Effect.flatMap((scope) =>
436 getPoolItemInner(self).pipe(
437 Effect.ensuring(Effect.sync(() => self.state.waiters--)),
438 Effect.tap((item) => {
439 if (item.exit._tag === "Failure") {
440 self.state.items.delete(item)
441 self.state.invalidated.delete(item)
442 self.state.available.delete(item)
443 return self.state.semaphore.release(1)
444 }
445 item.refCount++
446 self.state.available.delete(item)
447 if (item.refCount < self.config.concurrency) {
448 self.state.available.add(item)
449 }
450 return Scope.addFinalizerExit(scope, () =>
451 Effect.flatMap(
452 Effect.suspend(() => {
453 item.refCount--
454 if (self.state.invalidated.has(item)) {
455 return invalidatePoolItem(self, item)
456 }
457 self.state.available.add(item)
458 return Effect.void
459 }),
460 () => self.state.semaphore.release(1)
461 ))
462 }),
463 Effect.onInterrupt(() => self.state.semaphore.release(1))
464 )
465 )
466 )
467 )
468

Callers 1

getFunction · 0.85

Calls 8

invalidatePoolItemFunction · 0.85
syncMethod · 0.80
onInterruptMethod · 0.80
pipeMethod · 0.65
takeMethod · 0.65
releaseMethod · 0.65
addMethod · 0.65
hasMethod · 0.45

Tested by

no test coverage detected

Used in the wild real call sites across dependent graphs

searching dependent graphs…