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 020b253f5e..ac8ca621f3 100644 --- a/kernel/shared/src/main/scala/cats/effect/kernel/Resource.scala +++ b/kernel/shared/src/main/scala/cats/effect/kernel/Resource.scala @@ -774,6 +774,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 { @@ -1410,18 +1428,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, *]] { diff --git a/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala b/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala index 16f1b38a5f..f45bb0e90a 100644 --- a/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala +++ b/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala @@ -1170,6 +1170,29 @@ class ResourceSpec extends BaseSpec with ScalaCheck with Discipline { }.void must completeAs(()) } } + + // 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] + 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(identity) + .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 } "attempt" >> {