From 3cc297b9260fbf834a1cdff62d7ca288d2e22a7b Mon Sep 17 00:00:00 2001 From: Yannick Heiber Date: Wed, 9 Sep 2026 10:44:57 +0200 Subject: [PATCH 1/2] Add test cases demonstrating permit loss RequestSemaphore (and hence the pool) could lose permits when a waiting fiber gets cancelled with unfortunate timing. These test cases reproduce this behaviour to guide a fix and prevent regression. They're internally repeated 100 times as even with TestControl, there's some randomness as to which fiber is picked first when at a tick, several are ready. --- .../org/typelevel/keypool/PoolSpec.scala | 25 +++++++++++++++++++ .../internal/RequestSemaphoreSpec.scala | 25 +++++++++++++++++++ 2 files changed, 50 insertions(+) diff --git a/core/src/test/scala/org/typelevel/keypool/PoolSpec.scala b/core/src/test/scala/org/typelevel/keypool/PoolSpec.scala index 12f69552..e81367d2 100644 --- a/core/src/test/scala/org/typelevel/keypool/PoolSpec.scala +++ b/core/src/test/scala/org/typelevel/keypool/PoolSpec.scala @@ -26,6 +26,7 @@ import cats.effect._ import cats.effect.std.CountDownLatch import cats.effect.testkit.TestControl import scala.concurrent.duration._ +import scala.concurrent.TimeoutException import munit.CatsEffectSuite class PoolSpec extends CatsEffectSuite { @@ -257,6 +258,30 @@ class PoolSpec extends CatsEffectSuite { } } + test("do not lose permits when requests time out while another request releases") { + val program = Pool + .Builder(Ref.of[IO, Int](1), nothing) + .withMaxTotal(1) + .build + .use { pool => + for { + holder <- pool.take.use(_ => IO.sleep(1.second)).start + _ <- IO.sleep(1.milli) + // These timeouts fire at the same instant the holder releases. A request that is + // cancelled while `release` hands it the permit takes the permit down with it. + _ <- (1 to 8).toList.parTraverse_(_ => pool.take.use_.timeout(999.millis).attempt) + _ <- holder.join + // Test whether we can acquire a permit, fail with timeout if we can't (no permits available) + _ <- pool.take.use_.timeout(1.minute) + } yield () + } + + TestControl + .executeEmbed(program) + .adaptErr { case e: TimeoutException => new AssertionError(s"permit lost!", e) } + .replicateA_(100) // repeat to increase the chance of hitting the race condition + } + private def reqAction( pool: Pool[IO, Ref[IO, Int]], ref: Ref[IO, List[Int]], diff --git a/core/src/test/scala/org/typelevel/keypool/internal/RequestSemaphoreSpec.scala b/core/src/test/scala/org/typelevel/keypool/internal/RequestSemaphoreSpec.scala index e57a2641..79873241 100644 --- a/core/src/test/scala/org/typelevel/keypool/internal/RequestSemaphoreSpec.scala +++ b/core/src/test/scala/org/typelevel/keypool/internal/RequestSemaphoreSpec.scala @@ -22,6 +22,7 @@ package org.typelevel.keypool.internal import munit.CatsEffectSuite +import cats.syntax.all._ import cats.effect._ import cats.effect.testkit.TestControl import scala.concurrent.duration._ @@ -122,6 +123,30 @@ class RequestSemaphoreSpec extends CatsEffectSuite { } } + List(Fifo, Lifo).foreach { fairness => + test( + s"$fairness: do not lose a permit when a waiter is cancelled while being handed the permit" + ) { + val program = for { + sem <- RequestSemaphore[IO](fairness, 1) + held <- sem.permit.allocated + waiter <- sem.permit.surround(IO.unit).start + _ <- IO.sleep(1.milli) // the waiter is now queued for the permit + // Cancelling only schedules the waiter's cleanup. If `release` dequeues the waiter first, + // it hands the permit to a fiber that is already being cancelled and never releases it. + _ <- IO.both(held._2, waiter.cancel) + _ <- sem.permit.surround(IO.unit) // never completes if the permit was lost + } yield () + + TestControl + .executeEmbed(program) + .adaptErr { case e: TestControl.NonTerminationException => + new AssertionError(s"permit lost!", e) + } + .replicateA_(100) // repeat to increase the chance of hitting the race condition + } + } + private def action( sem: RequestSemaphore[IO], ref: Ref[IO, List[Int]], From 2d814e2d5d6aa2fd61514a83eb3ce75aea21adbb Mon Sep 17 00:00:00 2001 From: Yannick Heiber Date: Wed, 9 Sep 2026 11:05:03 +0200 Subject: [PATCH 2/2] Fix permit loss in RequestSemaphore Ports changes to the original cats-effect MiniSemaphore implementation to RequestSemaphore (commit in https://github.com/typelevel/cats-effect/commit/643ba3356627da1542fcef619107133c8c944fd0). This fixes a bug which could lead to permit loss (aka pool leaks) if a waiting fiber get cancelled, but still gets handed a permit which is then never returned and effectively lost. The two test cases added in the previous commit now happily pass. --- .../keypool/internal/RequestSemaphore.scala | 14 ++++++++------ 1 file changed, 8 insertions(+), 6 deletions(-) diff --git a/core/src/main/scala/org/typelevel/keypool/internal/RequestSemaphore.scala b/core/src/main/scala/org/typelevel/keypool/internal/RequestSemaphore.scala index a98badfb..e2cb5fbd 100644 --- a/core/src/main/scala/org/typelevel/keypool/internal/RequestSemaphore.scala +++ b/core/src/main/scala/org/typelevel/keypool/internal/RequestSemaphore.scala @@ -36,7 +36,7 @@ import org.typelevel.keypool.Fairness * the order in which requests acquire a permit. * * Derived from cats-effect MiniSemaphore - * https://github.com/typelevel/cats-effect/blob/v3.5.4/kernel/shared/src/main/scala/cats/effect/kernel/MiniSemaphore.scala#L29 + * https://github.com/typelevel/cats-effect/blob/v3.7.1/kernel/shared/src/main/scala/cats/effect/kernel/MiniSemaphore.scala#L29 */ private[keypool] abstract class RequestSemaphore[F[_]] { def permit: Resource[F, Unit] @@ -95,10 +95,10 @@ private[keypool] object RequestSemaphore { new RequestSemaphore[F] { private def acquire: F[Unit] = F.deferred[Unit].flatMap { wait => - val cleanup = state.update { case s @ State(waiting, permits) => - if (B.nonEmpty(waiting)) - State(B.cleanup(waiting, wait), permits) - else s + val cleanup = state.flatModify { case State(waiting, permits) => + State(B.cleanup(waiting, wait), permits) -> wait.complete(()).flatMap { won => + if (won) F.unit else release + } } state.flatModifyFull { case (poll, State(waiting, permits)) => @@ -113,7 +113,9 @@ private[keypool] object RequestSemaphore { state.flatModify { case State(waiting, permits) => if (B.nonEmpty(waiting)) { val (rest, next) = B.take(waiting) - State(rest, permits) -> next.complete(()).void + State(rest, permits) -> next.complete(()).flatMap { granted => + if (granted) F.unit else release + } } else State(waiting, permits + 1) -> F.unit }