Hyperlinkv0.8.0-beta.28

PartitionedSemaphore

PartitionedSemaphore.makeUnsafeconsteffect/PartitionedSemaphore.ts:115
<K = unknown>(options: {
  readonly permits: number
}): PartitionedSemaphore<K>

Constructs a PartitionedSemaphore synchronously, outside of Effect.

When to use

Use when you need to construct a partitioned semaphore synchronously outside an Effect workflow.

Details

Negative permit counts are clamped to 0. Non-finite permit counts create an unbounded semaphore whose acquire and release operations complete immediately.

constructorsmake
export const makeUnsafe = <K = unknown>(options: {
  readonly permits: number
}): PartitionedSemaphore<K> => {
  const maxPermits = Math.max(0, options.permits)

  if (!Number.isFinite(maxPermits)) {
    return {
      [PartitionedTypeId]: PartitionedTypeId,
      capacity: maxPermits,
      available: Effect.succeed(maxPermits),
      take: () => Effect.void,
      release: () => Effect.succeed(maxPermits),
      withPermits: () => (effect) => effect,
      withPermit: () => (effect) => effect,
      withPermitsIfAvailable: () => (effect) => Effect.asSome(effect)
    }
  }

  let totalPermits = maxPermits
  let waitingPermits = 0

  type Waiter = {
    permits: number
    readonly resume: () => void
  }

  const partitions = MutableHashMap.empty<K, Set<Waiter>>()
  let iterator = partitions[Symbol.iterator]()

  const releaseUnsafe = (permits: number): number => {
    while (permits > 0) {
      if (waitingPermits === 0) {
        totalPermits = Math.min(maxPermits, totalPermits + permits)
        return totalPermits
      }

      let state = iterator.next()
      if (state.done) {
        iterator = partitions[Symbol.iterator]()
        state = iterator.next()
        if (state.done) {
          return totalPermits
        }
      }

      const waiter = state.value[1].values().next().value
      if (waiter === undefined) {
        continue
      }

      waiter.permits -= 1
      waitingPermits -= 1

      if (waiter.permits === 0) {
        waiter.resume()
      }

      permits -= 1
    }

    return totalPermits
  }

  const take = (key: K, permits: number): Effect.Effect<void> => {
    if (permits <= 0) {
      return Effect.void
    }

    return Effect.callback<void>((resume) => {
      if (maxPermits < permits) {
        resume(Effect.never)
        return
      }

      if (totalPermits >= permits) {
        totalPermits -= permits
        resume(Effect.void)
        return
      }

      const needed = permits - totalPermits
      const taken = permits - needed
      if (totalPermits > 0) {
        totalPermits = 0
      }
      waitingPermits += needed

      const waiters = Option.getOrElse(
        MutableHashMap.get(partitions, key),
        () => {
          const set = new Set<Waiter>()
          MutableHashMap.set(partitions, key, set)
          return set
        }
      )

      const entry: Waiter = {
        permits: needed,
        resume: () => {
          cleanup()
          resume(Effect.void)
        }
      }

      const cleanup = () => {
        waiters.delete(entry)
        if (waiters.size === 0) {
          MutableHashMap.remove(partitions, key)
        }
      }

      waiters.add(entry)

      return Effect.sync(() => {
        cleanup()
        waitingPermits -= entry.permits
        if (taken > 0) {
          releaseUnsafe(taken)
        }
      })
    })
  }

  const withPermits =
    (key: K, permits: number) => <A, E, R>(effect: Effect.Effect<A, E, R>): Effect.Effect<A, E, R> => {
      if (permits <= 0) {
        return effect
      }

      const takePermits = take(key, permits)
      return Effect.uninterruptibleMask((restore) =>
        Effect.flatMap(
          restore(takePermits),
          () =>
            Effect.ensuring(
              restore(effect),
              Effect.sync(() => {
                releaseUnsafe(permits)
              })
            )
        )
      )
    }

  const tryTake = (permits: number): boolean => {
    if (permits <= 0) {
      return true
    }

    if (maxPermits < permits || totalPermits < permits) {
      return false
    }

    totalPermits -= permits
    return true
  }

  return {
    [PartitionedTypeId]: PartitionedTypeId,
    capacity: maxPermits,
    available: Effect.sync(() => totalPermits),
    take,
    release: (permits) => Effect.sync(() => releaseUnsafe(permits)),
    withPermits,
    withPermit: (key) => withPermits(key, 1),
    withPermitsIfAvailable:
      (permits) => <A, E, R>(effect: Effect.Effect<A, E, R>): Effect.Effect<Option.Option<A>, E, R> => {
        if (permits <= 0) {
          return Effect.asSome(effect)
        }

        return Effect.suspend(() => {
          if (!tryTake(permits)) {
            return Effect.succeed(Option.none())
          }

          return Effect.ensuring(
            Effect.asSome(effect),
            Effect.sync(() => {
              releaseUnsafe(permits)
            })
          )
        })
      }
  }
}
Referenced by 1 symbols