From 88d3a8587c131cdae4834675b3171c436b7035cc Mon Sep 17 00:00:00 2001 From: Daniel Urban Date: Sat, 27 Jun 2026 15:54:47 +0200 Subject: [PATCH 1/9] Fix #4489 --- .../cats/effect/kernel/GenConcurrent.scala | 9 ++-- .../scala/cats/effect/kernel/Resource.scala | 41 ++++++++++------ .../test/scala/cats/effect/ResourceSpec.scala | 49 +++++++++++++++++++ 3 files changed, 79 insertions(+), 20 deletions(-) diff --git a/kernel/shared/src/main/scala/cats/effect/kernel/GenConcurrent.scala b/kernel/shared/src/main/scala/cats/effect/kernel/GenConcurrent.scala index 024f640197..7ac122f7ec 100644 --- a/kernel/shared/src/main/scala/cats/effect/kernel/GenConcurrent.scala +++ b/kernel/shared/src/main/scala/cats/effect/kernel/GenConcurrent.scala @@ -176,10 +176,11 @@ trait GenConcurrent[F[_], E] extends GenSpawn[F, E] { _ <- canA.join _ <- canB.join } yield ()) - } yield back match { - case Left(oc) => Left((oc, fibB)) - case Right(oc) => Right((fibA, oc)) - } + result <- back match { + case Left(oc) => fibA.join.as(Left((oc, fibB))) + case Right(oc) => fibB.join.as(Right((fibA, oc))) + } + } yield result } } } diff --git a/kernel/shared/src/main/scala/cats/effect/kernel/Resource.scala b/kernel/shared/src/main/scala/cats/effect/kernel/Resource.scala index 7299b7afdb..f42be8e954 100644 --- a/kernel/shared/src/main/scala/cats/effect/kernel/Resource.scala +++ b/kernel/shared/src/main/scala/cats/effect/kernel/Resource.scala @@ -631,26 +631,35 @@ sealed abstract class Resource[F[_], +A] extends Serializable { F.ref[State](State()) flatMap { state => val finalized: F[A] = F uncancelable { poll => - poll(this.allocated) guarantee { + poll(this.allocated) guaranteeCase { // confirm that we completed and we were asked to clean up // note that this will run even if the inner effect short-circuited - state update { s => - if (s.finalizeOnComplete) - s.copy(confirmedFinalizeOnComplete = true) - else - s - } - } flatMap { - // if the inner F has a zero, we lose the finalizers, but there's no avoiding that - case (a, rel) => - val action = state modify { s => - if (s.confirmedFinalizeOnComplete) - (s, rel.handleError(_ => ())) + case Canceled() | Errored(_) => + state.update { s => + if (s.finalizeOnComplete) + s.copy(confirmedFinalizeOnComplete = true) else - (s.copy(fin = rel), F.unit) + s } - - action.flatten.as(a) + case Succeeded(fp) => + // if the inner F has a zero, we lose the finalizers, but there's no avoiding that + fp.flatMap { + case (_, rel) => + val action = state.modify { s => + if (s.finalizeOnComplete) { + // finalize immediately + (s.copy(confirmedFinalizeOnComplete = true), rel.voidError) + } else { + // save the finalizer for later + (s.copy(fin = rel), F.unit) + } + } + action.flatten + } + } map { + case (a, _) => + // Note: we've already saved/used the finalizer, see above + a } } diff --git a/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala b/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala index 625cd2a5fb..64ef341394 100644 --- a/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala +++ b/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala @@ -1193,6 +1193,55 @@ class ResourceSpec extends BaseSpec with ScalaCheck with Discipline { run(go) == Outcome.succeeded(Some(true)) }.pendingUntilFixed + + // issue #4489 repro (1) + "timeout finalizer (start/release race)" in real { + val test: IO[Unit] = IO.ref(false).flatMap { ref => + val res = Resource.make(ref.set(true))(_ => ref.set(false)) + val timedRes = res.timeout(1.hour) + timedRes.use_ *> ref.get.ifM(IO.raiseError(new Exception("not released")), IO.unit) + } + + test.replicateA_(1000).as(ok) + } + + // issue #4489 repro (2) + "racePair finalizer (start/release race)" in real { + val test: IO[Unit] = IO.ref(false).flatMap { ref => + val res = Resource.make(ref.set(true))(_ => ref.set(false)) + val racedRes = Spawn[Resource[IO, *]].racePair(res, Resource.never) + racedRes.use_ *> ref.get.ifM(IO.raiseError(new Exception("not released")), IO.unit) + } + + test.replicateA_(1000).as(ok) + } + + // issue #4489 repro (2) variant (never actually failed) + "racePair finalizer (start/cancel race)" in real { + val test: IO[Unit] = IO.ref(false).flatMap { ref => + val res = Resource.make(ref.set(true))(_ => ref.set(false)) + val racedRes = Spawn[Resource[IO, *]].racePair(res, Resource.unit).flatMap { + case Left((_, fib)) => fib.cancel + case Right((fib, _)) => fib.cancel + } + racedRes.use_ *> ref.get.ifM(IO.raiseError(new Exception("not released")), IO.unit) + } + + test.parReplicateA_(1000).as(ok) + } + + // issue #4489 repro (3) + "racePair finalizer variant (start/release race)" in real { + val test: IO[Unit] = IO.ref(false).flatMap { ref => + IO.deferred[Unit].flatMap { d => + val res = Resource.make(ref.set(true))(_ => ref.set(false) <* d.complete(())) + val racedRes = Spawn[Resource[IO, *]].racePair(res, Resource.never) + racedRes.use_ *> ref.get.ifM(d.get, IO.unit) // hangs if finalizer doesn't run + } + } + + test.replicateA_(1000).as(ok) + } } "attempt" >> { From 52cc0d4aa5b03b2b310169bb1ab85bb02875bb82 Mon Sep 17 00:00:00 2001 From: Daniel Urban Date: Sat, 27 Jun 2026 16:10:26 +0200 Subject: [PATCH 2/9] Make sure there is no #4059 regression --- .../test/scala/cats/effect/ResourceSpec.scala | 19 +++++++++++++++++++ 1 file changed, 19 insertions(+) diff --git a/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala b/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala index 64ef341394..f464c681df 100644 --- a/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala +++ b/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala @@ -1242,6 +1242,25 @@ class ResourceSpec extends BaseSpec with ScalaCheck with Discipline { test.replicateA_(1000).as(ok) } + + // issue #4059 test (1) specialized to Resource + "propagate successful result from a completed effect" in real { + Resource.catsEffectTemporalForResource[IO].sleep(50.millis).as(true).uncancelable.timeout(10.millis).use { res => + IO(res must beTrue) + } + } + + // issue #4059 test (2) specialized to Resource + "propagate error from a completed effect" in real { + Resource.catsEffectTemporalForResource[IO].sleep(50.millis).flatMap { _ => + Resource.raiseError[IO, Unit, Throwable](new RuntimeException) + }.uncancelable + .timeout(10.millis) + .attempt + .use { res => + IO(res must beLike { case Left(e) => e must haveClass[RuntimeException] }) + } + } } "attempt" >> { From c2ab728ac34bcd2be37d59f4e4bf3e5efffa8250 Mon Sep 17 00:00:00 2001 From: Daniel Urban Date: Sat, 27 Jun 2026 20:45:35 +0200 Subject: [PATCH 3/9] Plug another possible leak + add another test --- .../scala/cats/effect/kernel/Resource.scala | 4 ++-- .../test/scala/cats/effect/ResourceSpec.scala | 17 +++++++++++++++++ 2 files changed, 19 insertions(+), 2 deletions(-) diff --git a/kernel/shared/src/main/scala/cats/effect/kernel/Resource.scala b/kernel/shared/src/main/scala/cats/effect/kernel/Resource.scala index f42be8e954..264f566fe8 100644 --- a/kernel/shared/src/main/scala/cats/effect/kernel/Resource.scala +++ b/kernel/shared/src/main/scala/cats/effect/kernel/Resource.scala @@ -663,7 +663,7 @@ sealed abstract class Resource[F[_], +A] extends Serializable { } } - F.start(finalized) map { outer => + F.start(finalized).map { outer => val fiber = new Fiber[Resource[F, *], Throwable, A] { def cancel = Resource eval { @@ -697,7 +697,7 @@ sealed abstract class Resource[F[_], +A] extends Serializable { state.modify(s => (s.copy(finalizeOnComplete = true), s.fin)).flatten (fiber, finalizeOuter) - } + }.uncancelable } } } diff --git a/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala b/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala index f464c681df..f98bf50a04 100644 --- a/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala +++ b/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala @@ -1261,6 +1261,23 @@ class ResourceSpec extends BaseSpec with ScalaCheck with Discipline { IO(res must beLike { case Left(e) => e must haveClass[RuntimeException] }) } } + + // additional #4059 test for Resource + "timeout finalizer (#4059)" in real { + val test = IO.ref(false).flatMap { ref => + val program = Resource.make(ref.set(true) *> IO.sleep(10.millis)) { _ => + ref.set(false) + }.timeout(10.millis).use_ + program.attempt.flatMap { + case Left(_) => + ref.get.ifM(IO.raiseError(new Exception("not released")), IO.unit) + case Right(_) => + IO.unit // no timeout + } + } + + test.replicateA_(1000).as(ok) + } } "attempt" >> { From 0319feec3195cdaa91bdc9697cd4abbc6680c165 Mon Sep 17 00:00:00 2001 From: Daniel Urban Date: Sat, 27 Jun 2026 20:53:35 +0200 Subject: [PATCH 4/9] Formatttttttting --- .../scala/cats/effect/kernel/Resource.scala | 62 ++++++++++--------- .../test/scala/cats/effect/ResourceSpec.scala | 25 +++++--- 2 files changed, 48 insertions(+), 39 deletions(-) diff --git a/kernel/shared/src/main/scala/cats/effect/kernel/Resource.scala b/kernel/shared/src/main/scala/cats/effect/kernel/Resource.scala index 264f566fe8..cea2da3d4f 100644 --- a/kernel/shared/src/main/scala/cats/effect/kernel/Resource.scala +++ b/kernel/shared/src/main/scala/cats/effect/kernel/Resource.scala @@ -663,41 +663,43 @@ sealed abstract class Resource[F[_], +A] extends Serializable { } } - F.start(finalized).map { outer => - val fiber = new Fiber[Resource[F, *], Throwable, A] { - def cancel = - Resource eval { - F uncancelable { poll => - // technically cancel is uncancelable, but separation of concerns and what not - poll(outer.cancel) *> state.update(_.copy(finalizeOnComplete = true)) + F.start(finalized) + .map { outer => + val fiber = new Fiber[Resource[F, *], Throwable, A] { + def cancel = + Resource eval { + F uncancelable { poll => + // technically cancel is uncancelable, but separation of concerns and what not + poll(outer.cancel) *> state.update(_.copy(finalizeOnComplete = true)) + } } - } - def join = - Resource eval { - outer.join.flatMap[Outcome[Resource[F, *], Throwable, A]] { - case Canceled() => - Outcome.canceled[Resource[F, *], Throwable, A].pure[F] - - case Errored(e) => - Outcome.errored[Resource[F, *], Throwable, A](e).pure[F] - - case Succeeded(fp) => - state.get map { s => - if (s.confirmedFinalizeOnComplete) - Outcome.canceled[Resource[F, *], Throwable, A] - else - Outcome.succeeded(Resource.eval(fp)) - } + def join = + Resource eval { + outer.join.flatMap[Outcome[Resource[F, *], Throwable, A]] { + case Canceled() => + Outcome.canceled[Resource[F, *], Throwable, A].pure[F] + + case Errored(e) => + Outcome.errored[Resource[F, *], Throwable, A](e).pure[F] + + case Succeeded(fp) => + state.get map { s => + if (s.confirmedFinalizeOnComplete) + Outcome.canceled[Resource[F, *], Throwable, A] + else + Outcome.succeeded(Resource.eval(fp)) + } + } } - } - } + } - val finalizeOuter = - state.modify(s => (s.copy(finalizeOnComplete = true), s.fin)).flatten + val finalizeOuter = + state.modify(s => (s.copy(finalizeOnComplete = true), s.fin)).flatten - (fiber, finalizeOuter) - }.uncancelable + (fiber, finalizeOuter) + } + .uncancelable } } } diff --git a/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala b/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala index f98bf50a04..6efc4bf6b4 100644 --- a/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala +++ b/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala @@ -1245,16 +1245,22 @@ class ResourceSpec extends BaseSpec with ScalaCheck with Discipline { // issue #4059 test (1) specialized to Resource "propagate successful result from a completed effect" in real { - Resource.catsEffectTemporalForResource[IO].sleep(50.millis).as(true).uncancelable.timeout(10.millis).use { res => - IO(res must beTrue) - } + Resource + .catsEffectTemporalForResource[IO] + .sleep(50.millis) + .as(true) + .uncancelable + .timeout(10.millis) + .use { res => IO(res must beTrue) } } // issue #4059 test (2) specialized to Resource "propagate error from a completed effect" in real { - Resource.catsEffectTemporalForResource[IO].sleep(50.millis).flatMap { _ => - Resource.raiseError[IO, Unit, Throwable](new RuntimeException) - }.uncancelable + Resource + .catsEffectTemporalForResource[IO] + .sleep(50.millis) + .flatMap { _ => Resource.raiseError[IO, Unit, Throwable](new RuntimeException) } + .uncancelable .timeout(10.millis) .attempt .use { res => @@ -1265,9 +1271,10 @@ class ResourceSpec extends BaseSpec with ScalaCheck with Discipline { // additional #4059 test for Resource "timeout finalizer (#4059)" in real { val test = IO.ref(false).flatMap { ref => - val program = Resource.make(ref.set(true) *> IO.sleep(10.millis)) { _ => - ref.set(false) - }.timeout(10.millis).use_ + val program = Resource + .make(ref.set(true) *> IO.sleep(10.millis)) { _ => ref.set(false) } + .timeout(10.millis) + .use_ program.attempt.flatMap { case Left(_) => ref.get.ifM(IO.raiseError(new Exception("not released")), IO.unit) From f4a8565d94ac329e2dd4074682000eccf8d85f18 Mon Sep 17 00:00:00 2001 From: Daniel Urban Date: Sat, 27 Jun 2026 21:24:41 +0200 Subject: [PATCH 5/9] Seriously? --- tests/shared/src/test/scala/cats/effect/ResourceSpec.scala | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala b/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala index 6efc4bf6b4..93ce6ae93f 100644 --- a/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala +++ b/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala @@ -1248,7 +1248,7 @@ class ResourceSpec extends BaseSpec with ScalaCheck with Discipline { Resource .catsEffectTemporalForResource[IO] .sleep(50.millis) - .as(true) + .map(_ => true) .uncancelable .timeout(10.millis) .use { res => IO(res must beTrue) } From 7e677c9b64bef595f463452056c9c5b01bf69352 Mon Sep 17 00:00:00 2001 From: Daniel Urban Date: Sat, 27 Jun 2026 23:02:03 +0200 Subject: [PATCH 6/9] CI is slow --- tests/shared/src/test/scala/cats/effect/ResourceSpec.scala | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala b/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala index 93ce6ae93f..f4e3a6a01a 100644 --- a/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala +++ b/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala @@ -1283,7 +1283,7 @@ class ResourceSpec extends BaseSpec with ScalaCheck with Discipline { } } - test.replicateA_(1000).as(ok) + test.replicateA_(500).as(ok) } } From 48e456f2b86197fda19f1628649f709858b1cda2 Mon Sep 17 00:00:00 2001 From: Daniel Urban Date: Sat, 27 Jun 2026 23:22:03 +0200 Subject: [PATCH 7/9] Let's try it in parallel --- tests/shared/src/test/scala/cats/effect/ResourceSpec.scala | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala b/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala index f4e3a6a01a..5d1f21652f 100644 --- a/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala +++ b/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala @@ -1283,7 +1283,7 @@ class ResourceSpec extends BaseSpec with ScalaCheck with Discipline { } } - test.replicateA_(500).as(ok) + test.parReplicateA_(1000).as(ok) } } From aa6101921bb1a14c7f3d61591c4647e2ffc2d707 Mon Sep 17 00:00:00 2001 From: Daniel Urban Date: Mon, 29 Jun 2026 10:44:42 +0200 Subject: [PATCH 8/9] Cleanup --- .../src/main/scala/cats/effect/kernel/Resource.scala | 11 ++++++----- 1 file changed, 6 insertions(+), 5 deletions(-) diff --git a/kernel/shared/src/main/scala/cats/effect/kernel/Resource.scala b/kernel/shared/src/main/scala/cats/effect/kernel/Resource.scala index cea2da3d4f..05479f74c6 100644 --- a/kernel/shared/src/main/scala/cats/effect/kernel/Resource.scala +++ b/kernel/shared/src/main/scala/cats/effect/kernel/Resource.scala @@ -621,15 +621,11 @@ sealed abstract class Resource[F[_], +A] extends Serializable { def start( implicit F: Concurrent[F]): Resource[F, Fiber[Resource[F, *], Throwable, A @uncheckedVariance]] = { - final case class State( - fin: F[Unit] = F.unit, - finalizeOnComplete: Boolean = false, - confirmedFinalizeOnComplete: Boolean = false) Resource { import Outcome._ - F.ref[State](State()) flatMap { state => + F.ref[Resource.FiberState[F]](Resource.FiberState(F.unit)) flatMap { state => val finalized: F[A] = F uncancelable { poll => poll(this.allocated) guaranteeCase { // confirm that we completed and we were asked to clean up @@ -807,6 +803,11 @@ sealed abstract class Resource[F[_], +A] extends Serializable { object Resource extends ResourceFOInstances0 with ResourceHOInstances0 with ResourcePlatform { + private final case class FiberState[F[_]]( + fin: F[Unit], + finalizeOnComplete: Boolean = false, + confirmedFinalizeOnComplete: Boolean = false) + /** * Creates a resource from an allocating effect. * From 3b7ddb059a9ad188c2176d8558b216c5aab69fc5 Mon Sep 17 00:00:00 2001 From: Daniel Urban Date: Mon, 29 Jun 2026 21:44:45 +0200 Subject: [PATCH 9/9] Improved test --- tests/shared/src/test/scala/cats/effect/ResourceSpec.scala | 7 ++----- 1 file changed, 2 insertions(+), 5 deletions(-) diff --git a/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala b/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala index 5d1f21652f..01e42b511d 100644 --- a/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala +++ b/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala @@ -1275,11 +1275,8 @@ class ResourceSpec extends BaseSpec with ScalaCheck with Discipline { .make(ref.set(true) *> IO.sleep(10.millis)) { _ => ref.set(false) } .timeout(10.millis) .use_ - program.attempt.flatMap { - case Left(_) => - ref.get.ifM(IO.raiseError(new Exception("not released")), IO.unit) - case Right(_) => - IO.unit // no timeout + program.attempt.flatMap { _ => + ref.get.ifM(IO.raiseError(new Exception("not released")), IO.unit) } }