From d04231edd7dec3d49d8ce19ac3277e41e385d0e6 Mon Sep 17 00:00:00 2001 From: Gavin Bisesi Date: Wed, 8 Nov 2023 10:17:57 -0500 Subject: [PATCH 1/8] Add a memoizedAcquire method to Resource Ref #3513 --- .../scala/cats/effect/kernel/Resource.scala | 18 ++++++++++++++++++ 1 file changed, 18 insertions(+) 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 0017996688..c503d33377 100644 --- a/kernel/shared/src/main/scala/cats/effect/kernel/Resource.scala +++ b/kernel/shared/src/main/scala/cats/effect/kernel/Resource.scala @@ -811,6 +811,24 @@ object Resource extends ResourceFOInstances0 with ResourceHOInstances0 with Reso } } + implicit final class IdSyntax[F[_], A](val self: Resource[F, A]) extends AnyVal { + + /** + * A Resource where the acquire step is done lazily and memoized. + * This means that acquire happens only if and when the `F[A]` value is executed, instead of happening immediately upon `use()`. + * If the `F[A]` value is executed multiple times, acquire happens once only and the acquired resource is shared to all callers. + * The resource is released as normal at the end of `use` (whether normal termination, error, or cancelled), if it was acquired. + */ + def memoizedAcquire(implicit F: Concurrent[F]): Resource[F, F[A]] = + // NB this uses syntax instead of an instance method because of variance on `A` within the class + Concurrent[Resource[F, *]].memoize[A](self).map { + case Resource.Eval(fa) => fa + case unexpected => + throw new IllegalStateException( + s"Memoized Resource is not Resource.Eval: $unexpected") + } + } + type Par[F[_], A] = ParallelF[Resource[F, *], A] implicit def parallelForResource[F[_]: Concurrent]: Parallel.Aux[Resource[F, *], Par[F, *]] = From e44f74bfebb05cf8b8e32dba45483f7838127e1f Mon Sep 17 00:00:00 2001 From: Arman Bilge Date: Mon, 20 May 2024 20:42:55 +0000 Subject: [PATCH 2/8] Refactor --- .../scala/cats/effect/kernel/Resource.scala | 50 ++++++++----------- 1 file changed, 20 insertions(+), 30 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 ff1a6284c5..2e14c90526 100644 --- a/kernel/shared/src/main/scala/cats/effect/kernel/Resource.scala +++ b/kernel/shared/src/main/scala/cats/effect/kernel/Resource.scala @@ -751,6 +751,24 @@ sealed abstract class Resource[F[_], +A] extends Serializable { K.combineK(allocate(this), allocate(that)) } + /** + * A Resource where the acquire step is done lazily and memoized. This means that acquire + * happens only if and when the `F[A]` value is executed, instead of happening immediately + * upon `use()`. If the `F[A]` value is executed multiple times, acquire happens once only and + * the acquired resource is shared to all callers. The resource is released as normal at the + * end of `use` (whether normal termination, error, or cancelled), if it was acquired. + */ + def memoizedAcquire[B >: A](implicit F: Concurrent[F]): Resource[F, F[B]] = { + Resource.eval(F.ref(List.empty[Resource.ExitCase => F[Unit]])).flatMap { release => + val fa2 = F.uncancelable { poll => + poll(allocatedCase).flatMap { case (a, r) => release.update(r :: _).as(a) } + } + Resource.makeCaseFull[F, F[B]](poll => poll(F.memoize(fa2).map(_.widen))) { (_, exit) => + release.get.flatMap(_.foldMapM(_(exit))) + } + } + } + } object Resource extends ResourceFOInstances0 with ResourceHOInstances0 with ResourcePlatform { @@ -1209,24 +1227,6 @@ object Resource extends ResourceFOInstances0 with ResourceHOInstances0 with Reso } } - implicit final class IdSyntax[F[_], A](val self: Resource[F, A]) extends AnyVal { - - /** - * A Resource where the acquire step is done lazily and memoized. - * This means that acquire happens only if and when the `F[A]` value is executed, instead of happening immediately upon `use()`. - * If the `F[A]` value is executed multiple times, acquire happens once only and the acquired resource is shared to all callers. - * The resource is released as normal at the end of `use` (whether normal termination, error, or cancelled), if it was acquired. - */ - def memoizedAcquire(implicit F: Concurrent[F]): Resource[F, F[A]] = - // NB this uses syntax instead of an instance method because of variance on `A` within the class - Concurrent[Resource[F, *]].memoize[A](self).map { - case Resource.Eval(fa) => fa - case unexpected => - throw new IllegalStateException( - s"Memoized Resource is not Resource.Eval: $unexpected") - } - } - type Par[F[_], A] = ParallelF[Resource[F, *], A] implicit def parallelForResource[F[_]: Concurrent]: Parallel.Aux[Resource[F, *], Par[F, *]] = @@ -1395,18 +1395,8 @@ abstract private[effect] class ResourceConcurrent[F[_]] override def race[A, B](fa: Resource[F, A], fb: Resource[F, B]): Resource[F, Either[A, B]] = fa.race(fb) - override def memoize[A](fa: Resource[F, A]): Resource[F, Resource[F, A]] = { - Resource.eval(F.ref(List.empty[Resource.ExitCase => F[Unit]])).flatMap { release => - val fa2 = F.uncancelable { poll => - poll(fa.allocatedCase).flatMap { case (a, r) => release.update(r :: _).as(a) } - } - Resource - .makeCaseFull[F, F[A]](poll => poll(F.memoize(fa2))) { (_, exit) => - release.get.flatMap(_.foldMapM(_(exit))) - } - .map(memo => Resource.eval(memo)) - } - } + override def memoize[A](fa: Resource[F, A]): Resource[F, Resource[F, A]] = + fa.memoizedAcquire.map(Resource.eval(_)) } private[effect] trait ResourceClock[F[_]] extends Clock[Resource[F, *]] { From 41dff47edb552971eb86040c89475f5a0b20d585 Mon Sep 17 00:00:00 2001 From: Lucas Satabin Date: Tue, 16 Jul 2024 21:02:53 +0200 Subject: [PATCH 3/8] Ensure acquire is not canceled once memoization succeeded --- kernel/shared/src/main/scala/cats/effect/kernel/Resource.scala | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 2e14c90526..3ace2ad2a0 100644 --- a/kernel/shared/src/main/scala/cats/effect/kernel/Resource.scala +++ b/kernel/shared/src/main/scala/cats/effect/kernel/Resource.scala @@ -763,7 +763,7 @@ sealed abstract class Resource[F[_], +A] extends Serializable { val fa2 = F.uncancelable { poll => poll(allocatedCase).flatMap { case (a, r) => release.update(r :: _).as(a) } } - Resource.makeCaseFull[F, F[B]](poll => poll(F.memoize(fa2).map(_.widen))) { (_, exit) => + Resource.makeCaseFull[F, F[B]](poll => poll(F.memoize(fa2)).map(_.widen)) { (_, exit) => release.get.flatMap(_.foldMapM(_(exit))) } } From 0626c7cbb05dff2b49dd1a903fa275e551d5f1e5 Mon Sep 17 00:00:00 2001 From: Lucas Satabin Date: Thu, 25 Jul 2024 18:29:41 +0200 Subject: [PATCH 4/8] Attempt at testing the cancellation --- .../test/scala/cats/effect/ResourceSpec.scala | 27 +++++++++++++++++++ 1 file changed, 27 insertions(+) diff --git a/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala b/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala index b917818531..eb6c663be1 100644 --- a/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala +++ b/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala @@ -36,6 +36,7 @@ import scala.concurrent.ExecutionContext import scala.concurrent.duration._ import java.util.concurrent.atomic.AtomicBoolean +import cats.effect.unsafe.IORuntimeConfig class ResourceSpec extends BaseSpec with ScalaCheck with Discipline { // We need this for testing laws: prop runs can interfere with each other @@ -1170,6 +1171,32 @@ class ResourceSpec extends BaseSpec with ScalaCheck with Discipline { }.void must completeAs(()) } } + + "does not leak if canceled right after delayed acquire is canceled" in ticked { implicit ticker => + (IO(new AtomicBoolean), IO.ref(false), IO.ref(false)).flatMapN { (acquired, used, released) => + val go = for { + fiber <- Resource + .eval(IO(acquired.set(true))) + .memoizedAcquire + .use(_ *> used.set(true)) + .start + _ <- IO.cede.untilM_(IO(acquired.get)) + _ <- fiber.cancel + _ <- fiber.join + } yield () + TestControl + .executeEmbed(go, IORuntimeConfig(1, 2)) + .flatMap { _ => + for { + acquireRun <- IO(acquired.get) + useRun <- used.get + releaseRun <- released.get + } yield acquireRun && releaseRun + } + .replicateA(1000) + .map(_.forall(identity(_))) + } must completeAs(true) + } } "uncancelable" >> { From ca5f5003681477e2a02a6a60e6e0e58d8f803b95 Mon Sep 17 00:00:00 2001 From: Lucas Satabin Date: Thu, 25 Jul 2024 18:35:31 +0200 Subject: [PATCH 5/8] Actually release something --- .../scala/cats/effect/kernel/Resource.scala | 2 +- .../test/scala/cats/effect/ResourceSpec.scala | 49 ++++++++++--------- 2 files changed, 26 insertions(+), 25 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 3ace2ad2a0..2e14c90526 100644 --- a/kernel/shared/src/main/scala/cats/effect/kernel/Resource.scala +++ b/kernel/shared/src/main/scala/cats/effect/kernel/Resource.scala @@ -763,7 +763,7 @@ sealed abstract class Resource[F[_], +A] extends Serializable { val fa2 = F.uncancelable { poll => poll(allocatedCase).flatMap { case (a, r) => release.update(r :: _).as(a) } } - Resource.makeCaseFull[F, F[B]](poll => poll(F.memoize(fa2)).map(_.widen)) { (_, exit) => + Resource.makeCaseFull[F, F[B]](poll => poll(F.memoize(fa2).map(_.widen))) { (_, exit) => release.get.flatMap(_.foldMapM(_(exit))) } } diff --git a/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala b/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala index eb6c663be1..e12bbdd9ac 100644 --- a/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala +++ b/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala @@ -1172,30 +1172,31 @@ class ResourceSpec extends BaseSpec with ScalaCheck with Discipline { } } - "does not leak if canceled right after delayed acquire is canceled" in ticked { implicit ticker => - (IO(new AtomicBoolean), IO.ref(false), IO.ref(false)).flatMapN { (acquired, used, released) => - val go = for { - fiber <- Resource - .eval(IO(acquired.set(true))) - .memoizedAcquire - .use(_ *> used.set(true)) - .start - _ <- IO.cede.untilM_(IO(acquired.get)) - _ <- fiber.cancel - _ <- fiber.join - } yield () - TestControl - .executeEmbed(go, IORuntimeConfig(1, 2)) - .flatMap { _ => - for { - acquireRun <- IO(acquired.get) - useRun <- used.get - releaseRun <- released.get - } yield acquireRun && releaseRun - } - .replicateA(1000) - .map(_.forall(identity(_))) - } must completeAs(true) + "does not leak if canceled right after delayed acquire is canceled" in ticked { + implicit ticker => + (IO(new AtomicBoolean), IO.ref(false), IO.ref(false)).flatMapN { + (acquired, used, released) => + val go = for { + fiber <- Resource + .make(IO(acquired.set(true)))(_ => released.set(true)) + .memoizedAcquire + .use(_ *> used.set(true)) + .start + _ <- IO.cede.untilM_(IO(acquired.get)) + _ <- fiber.cancel + } yield () + TestControl + .executeEmbed(go, IORuntimeConfig(1, 2)) + .flatMap { _ => + for { + acquireRun <- IO(acquired.get) + useRun <- used.get + releaseRun <- released.get + } yield acquireRun && releaseRun + } + .replicateA(1000) + .map(_.forall(identity(_))) + } must completeAs(true) } } From fdbf8cf9192bb27ac225fb0aafc67cbeac7994e8 Mon Sep 17 00:00:00 2001 From: Lucas Satabin Date: Sun, 1 Dec 2024 11:24:04 +0100 Subject: [PATCH 6/8] Add test for checking no leak occurs if cancelled --- .../scala/cats/effect/kernel/Resource.scala | 2 +- .../test/scala/cats/effect/ResourceSpec.scala | 51 +++++++++---------- 2 files changed, 25 insertions(+), 28 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 2e14c90526..3ace2ad2a0 100644 --- a/kernel/shared/src/main/scala/cats/effect/kernel/Resource.scala +++ b/kernel/shared/src/main/scala/cats/effect/kernel/Resource.scala @@ -763,7 +763,7 @@ sealed abstract class Resource[F[_], +A] extends Serializable { val fa2 = F.uncancelable { poll => poll(allocatedCase).flatMap { case (a, r) => release.update(r :: _).as(a) } } - Resource.makeCaseFull[F, F[B]](poll => poll(F.memoize(fa2).map(_.widen))) { (_, exit) => + Resource.makeCaseFull[F, F[B]](poll => poll(F.memoize(fa2)).map(_.widen)) { (_, exit) => release.get.flatMap(_.foldMapM(_(exit))) } } diff --git a/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala b/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala index e12bbdd9ac..8f0b560c36 100644 --- a/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala +++ b/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala @@ -36,7 +36,6 @@ import scala.concurrent.ExecutionContext import scala.concurrent.duration._ import java.util.concurrent.atomic.AtomicBoolean -import cats.effect.unsafe.IORuntimeConfig class ResourceSpec extends BaseSpec with ScalaCheck with Discipline { // We need this for testing laws: prop runs can interfere with each other @@ -1172,32 +1171,30 @@ class ResourceSpec extends BaseSpec with ScalaCheck with Discipline { } } - "does not leak if canceled right after delayed acquire is canceled" in ticked { - implicit ticker => - (IO(new AtomicBoolean), IO.ref(false), IO.ref(false)).flatMapN { - (acquired, used, released) => - val go = for { - fiber <- Resource - .make(IO(acquired.set(true)))(_ => released.set(true)) - .memoizedAcquire - .use(_ *> used.set(true)) - .start - _ <- IO.cede.untilM_(IO(acquired.get)) - _ <- fiber.cancel - } yield () - TestControl - .executeEmbed(go, IORuntimeConfig(1, 2)) - .flatMap { _ => - for { - acquireRun <- IO(acquired.get) - useRun <- used.get - releaseRun <- released.get - } yield acquireRun && releaseRun - } - .replicateA(1000) - .map(_.forall(identity(_))) - } must completeAs(true) - } + // TODO enable once `PureConc` finalizer bug is fixed. Test is indefinitely waiting for finalizers as of now + /* + "does not leak if canceled right after delayed acquire is canceled" in { + import cats.effect.kernel.testkit.pure._ + type F[A] = PureConc[Throwable, A] + val F = Concurrent[F] + def go = for { + acquired <- F.ref(false) + released <- F.ref(false) + fiber <- Resource + .make(acquired.set(true))(_ => released.set(true)) + .memoizedAcquire + .use_ + .start + _ <- F.cede.untilM_(acquired.get) + _ <- fiber.cancel + _ <- fiber.join + acquireRun <- acquired.get + releaseRun <- released.get + } yield acquireRun && releaseRun + + run(go) == Outcome.succeeded(Some(true)) + }.pendingUntilFixed + */ } "uncancelable" >> { From 7ed1e28796c00934583ad2d499a267545cf22c0c Mon Sep 17 00:00:00 2001 From: Lucas Satabin Date: Sun, 1 Dec 2024 12:58:30 +0100 Subject: [PATCH 7/8] Run `prePR` --- 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 8439a17a2f..9fe18cf9d9 100644 --- a/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala +++ b/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala @@ -1194,7 +1194,7 @@ class ResourceSpec extends BaseSpec with ScalaCheck with Discipline { run(go) == Outcome.succeeded(Some(true)) }.pendingUntilFixed - */ + */ } "attempt" >> { From cd6103377f4abe93a569ae8da52b974bcfd257d4 Mon Sep 17 00:00:00 2001 From: Lucas Satabin Date: Sun, 1 Dec 2024 16:08:24 +0100 Subject: [PATCH 8/8] Mark test as pending --- tests/shared/src/test/scala/cats/effect/ResourceSpec.scala | 6 ++---- 1 file changed, 2 insertions(+), 4 deletions(-) diff --git a/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala b/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala index 9fe18cf9d9..f45bb0e90a 100644 --- a/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala +++ b/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala @@ -1171,8 +1171,7 @@ class ResourceSpec extends BaseSpec with ScalaCheck with Discipline { } } - // TODO enable once `PureConc` finalizer bug is fixed. Test is indefinitely waiting for finalizers as of now - /* + // TODO enable once `PureConc` finalizer bug is fixed. "does not leak if canceled right after delayed acquire is canceled" in { import cats.effect.kernel.testkit.pure._ type F[A] = PureConc[Throwable, A] @@ -1183,7 +1182,7 @@ class ResourceSpec extends BaseSpec with ScalaCheck with Discipline { fiber <- Resource .make(acquired.set(true))(_ => released.set(true)) .memoizedAcquire - .use_ + .use(identity) .start _ <- F.cede.untilM_(acquired.get) _ <- fiber.cancel @@ -1194,7 +1193,6 @@ class ResourceSpec extends BaseSpec with ScalaCheck with Discipline { run(go) == Outcome.succeeded(Some(true)) }.pendingUntilFixed - */ } "attempt" >> {