Skip to content
Open
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
17 changes: 11 additions & 6 deletions core/shared/src/main/scala/cats/effect/IOFiber.scala
Original file line number Diff line number Diff line change
Expand Up @@ -866,12 +866,14 @@ private final class IOFiber[A](

if (!shouldFinalize()) {
/* we weren't canceled, so resume the runloop */
val next = result match {
case Left(t) => failed(t, 0)
case Right(a) => succeeded(a, 0)
result match {
case Left(t) if !UnsafeNonFatal(t) =>
onFatalFailure(t)
case Left(t) =>
runLoop(failed(t, 0), nextCancelation, nextAutoCede)
case Right(a) =>
runLoop(succeeded(a, 0), nextCancelation, nextAutoCede)
}

runLoop(next, nextCancelation, nextAutoCede)
} else if (outcome == null) {
/*
* we were canceled, but `cancel` cannot run the finalisers
Expand Down Expand Up @@ -1407,7 +1409,10 @@ private final class IOFiber[A](

private[this] def asyncContinueFailedR(): Unit = {
val t = objectState.pop().asInstanceOf[Throwable]
runLoop(failed(t, 0), runtime.cancelationCheckThreshold, runtime.autoYieldThreshold)
if (!UnsafeNonFatal(t))
onFatalFailure(t)
else
runLoop(failed(t, 0), runtime.cancelationCheckThreshold, runtime.autoYieldThreshold)
}

private[this] def asyncContinueCanceledR(): Unit = {
Expand Down
15 changes: 15 additions & 0 deletions ioapp-tests/src/test/scala/IOAppSuite.scala
Original file line number Diff line number Diff line change
Expand Up @@ -225,6 +225,13 @@ class IOAppSuite extends FunSuite {
assert(!h.stdout().contains("sadness"))
}

test("exit on a fatal error surfaced through async_ with attempt") {
val h = platform("AsyncFatalError", List.empty)
assertEquals(h.awaitStatus(), 1)
assert(h.stderr().contains("Boom!"))
assert(!h.stdout().contains("sadness"))
}

test("warn on global runtime collision") {
val h = platform("GlobalRacingInit", List.empty)
assertEquals(h.awaitStatus(), 0)
Expand Down Expand Up @@ -386,6 +393,14 @@ class IOAppSuite extends FunSuite {
assert(!h.stdout().contains("sadness"))
assert(h.stdout().contains("done"))
}

test(
"exit on a fatal error surfaced through a fiber started via `fromCompletableFuture`") {
val h = platform("FatalErrorFromCompletableFuture", List.empty)
assertEquals(h.awaitStatus(), 1)
assert(h.stderr().contains("Boom!"))
assert(!h.stdout().contains("sadness"))
}
}

if (platform == Node) {
Expand Down
1 change: 1 addition & 0 deletions tests/js/src/main/scala/catseffect/examplesplatform.scala
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,7 @@ package examples {
register(RaiseFatalErrorHandle)
register(RaiseFatalErrorMap)
register(RaiseFatalErrorFlatMap)
register(AsyncFatalError)
registerRaw(FatalErrorRaw)
register(Canceled)
registerLazy("catseffect.examples.GlobalRacingInit", GlobalRacingInit)
Expand Down
33 changes: 32 additions & 1 deletion tests/jvm/src/main/scala/catseffect/examplesplatform.scala
Original file line number Diff line number Diff line change
Expand Up @@ -16,12 +16,13 @@

package catseffect

import cats.effect.{ExitCode, IO, IOApp}
import cats.effect.{ExitCode, IO, IOApp, Resource}
import cats.syntax.all._

import scala.concurrent.ExecutionContext
import scala.concurrent.duration._

import java.util.concurrent.{CompletableFuture, Executor, Executors}
import java.util.concurrent.atomic.AtomicReference

package object examples {
Expand Down Expand Up @@ -95,4 +96,34 @@ package examples {
val run =
IO.cede.foreverM.start >> IO(Thread.sleep(2.seconds.toMillis))
}

object FatalErrorFromCompletableFuture extends IOApp {

private val pingIO =
(IO.println("ping") *> IO.sleep(1.seconds)).foreverM

private def boomFromCompletableFuture(executor: Executor): IO[Unit] =
IO.fromCompletableFuture(
IO(
CompletableFuture.runAsync(
() => {
println("Waiting 2 seconds before boom...")
Thread.sleep(2000)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Is there a reason to sleep?

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

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

I just copied the example verbatim from the original ticket, the wait is there: #4505

println("Gonna boom!")
throw new OutOfMemoryError("Boom!")
},
executor)))
.void

override def run(args: List[String]): IO[ExitCode] =
Resource.make(IO(Executors.newFixedThreadPool(1)))(es => IO(es.shutdown())).use {
executor =>
for {
pingFiber <- pingIO.start
_ <- boomFromCompletableFuture(ExecutionContext.fromExecutor(executor)).start
_ <- pingFiber.join
} yield ExitCode.Success
}

}
}
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,7 @@ package examples {
register(RaiseFatalErrorHandle)
register(RaiseFatalErrorMap)
register(RaiseFatalErrorFlatMap)
register(AsyncFatalError)
registerRaw(FatalErrorRaw)
register(Canceled)
registerLazy("catseffect.examples.GlobalRacingInit", GlobalRacingInit)
Expand Down
14 changes: 12 additions & 2 deletions tests/shared/src/main/scala/catseffect/examples.scala
Original file line number Diff line number Diff line change
Expand Up @@ -53,7 +53,7 @@ package examples {

object FatalErrorRaw extends RawApp {
def main(args: Array[String]): Unit = {
import cats.effect.unsafe.implicits._
import cats.effect.unsafe.implicits.*
val action =
IO(throw new OutOfMemoryError("Boom!")).attempt.flatMap(_ => IO.println("sadness"))
action.unsafeToFuture()
Expand Down Expand Up @@ -131,6 +131,16 @@ package examples {
}
}

object AsyncFatalError extends IOApp {
def run(args: List[String]): IO[ExitCode] = {
IO.async_[Unit](cb => cb(Left(new OutOfMemoryError("Boom!"))))
.start
.flatMap(_.join)
.flatMap(_ => IO.println("sadness"))

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

The original reproducer in #4505 runs the async on a fiber and joins it:

        _ <- pingFiber.join
      } yield ExitCode.Success

join returns an object with the outcome of the fiber. This ignores the actual outcome of the async call. Your implementation doesn't use a fiber, so the failure is still propagated and causes your IO to fail. pingFiber also fails, the issue is it doesn't fatally fail.

The core issue with #4505 is that the code shouldn't even get far enough that there is a join result, as the OOM should kill the app immediately. You can probably reproduce this with just an attempt instead of a start/join

Suggested change
.flatMap(_ => IO.println("sadness"))
.attempt
.flatMap(_ => IO.println("sadness"))

or to more closely replicate the original, something like:

Suggested change
.flatMap(_ => IO.println("sadness"))
.start
.flatMap(_.join)
.flatMap(_ => IO.println("sadness"))

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

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

Applied the second change here, it made the test red: ec8421b

Thanks!

.as(ExitCode.Success)
}
}

object Canceled extends IOApp {
def run(args: List[String]): IO[ExitCode] =
IO.canceled.as(ExitCode.Success)
Expand Down Expand Up @@ -171,7 +181,7 @@ package examples {

object LiveFiberSnapshot extends IOApp.Simple {

import scala.concurrent.duration._
import scala.concurrent.duration.*

lazy val loop: IO[Unit] =
IO.unit.map(_ => ()) >>
Expand Down
Loading