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..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,74 +621,81 @@ 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) 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 } } - 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) - } + (fiber, finalizeOuter) + } + .uncancelable } } } @@ -796,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. * diff --git a/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala b/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala index 625cd2a5fb..01e42b511d 100644 --- a/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala +++ b/tests/shared/src/test/scala/cats/effect/ResourceSpec.scala @@ -1193,6 +1193,95 @@ 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) + } + + // issue #4059 test (1) specialized to Resource + "propagate successful result from a completed effect" in real { + Resource + .catsEffectTemporalForResource[IO] + .sleep(50.millis) + .map(_ => 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] }) + } + } + + // 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 { _ => + ref.get.ifM(IO.raiseError(new Exception("not released")), IO.unit) + } + } + + test.parReplicateA_(1000).as(ok) + } } "attempt" >> {