Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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
}
}
}
Expand Down
114 changes: 63 additions & 51 deletions kernel/shared/src/main/scala/cats/effect/kernel/Resource.scala
Original file line number Diff line number Diff line change
Expand Up @@ -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(

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This is an unrelated change, but local case classes have a strange encoding with a LazyRef, so I've taken the opportunity to clean it up. (It's moved to the companion object.)

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
}
}
}
Expand Down Expand Up @@ -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.
*
Expand Down
89 changes: 89 additions & 0 deletions tests/shared/src/test/scala/cats/effect/ResourceSpec.scala
Original file line number Diff line number Diff line change
Expand Up @@ -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" >> {
Expand Down
Loading