Skip to content

Queue.take.timeout(...) loses elements on cancelation #4571

Description

@TomasMikula

While queue.take.timeout(...) by itself does not lose elements on timeout, it does lose elements when the fiber running it is canceled.

Minimized code

import cats.effect.std.Queue
import cats.effect.{ExitCode, IO, IOApp}

import scala.concurrent.TimeoutException
import scala.concurrent.duration.DurationInt

object QueueTakeTimeoutCancelTest1 extends IOApp {
  val delay = 1.microsecond

  override def run(args: List[String]): IO[ExitCode] =
    Queue
      .unbounded[IO, Int]
      .flatMap { queue =>
        val check: IO[Unit] =
          for {
            // producer will enqueue an element after the given `delay`
            _ <- (IO.sleep(delay) *> queue.offer(42)).start

            // consumer will wait for an elem for up to the given `delay`
            elem <- queue.take
              .timeout(delay)
              .timeout(delay) // THIS ONE CAUSES TROUBLE (without this line, the test passes)
              .recover { case _: TimeoutException => -1 }

            _ <- IO.whenA(elem == -1) {
              // Timed out, the element must not have been dequeued.
              // 1 second should be plenty of time to dequeue it (TimeoutException means element was lost)
              queue.take.timeout(1.second)
                .adaptError { case _: TimeoutException => new RuntimeException("Element was lost!") }
                .void
            }
          } yield ()

        check.replicateA_(10000)
      }
      .as(ExitCode.Success)
}

Output

java.lang.RuntimeException: Element was lost!
	at QueueTakeTimeoutCancelTest1$$anonfun$$nestedInanonfun$run$4$1.applyOrElse(QueueTakeTimeoutCancelTest1.scala:29)
	at QueueTakeTimeoutCancelTest1$$anonfun$$nestedInanonfun$run$4$1.applyOrElse(QueueTakeTimeoutCancelTest1.scala:29)
	at scala.PartialFunction$AndThen.applyOrElse(PartialFunction.scala:299)
	at timeout @ QueueTakeTimeoutCancelTest1$.$anonfun$run$2(QueueTakeTimeoutCancelTest1.scala:21)
	at main$ @ QueueTakeTimeoutCancelTest1$.main(QueueTakeTimeoutCancelTest1.scala:7)
	at main$ @ QueueTakeTimeoutCancelTest1$.main(QueueTakeTimeoutCancelTest1.scala:7)

More realistic example

This example uses explicit cancelation (instead of nested timeouts) of a consumer fiber that's calling queue.take.timeout(...).

This mimics my actual use case, namely a consumer that:

  • blocks for a limited time when attempting to dequeue an element
  • never loses a dequeued element
  • can itself be canceled from outside, incl. when waiting on the queue
import cats.effect.kernel.{Deferred, Outcome}
import cats.effect.std.Queue
import cats.effect.{ExitCode, IO, IOApp}

import scala.concurrent.TimeoutException
import scala.concurrent.duration.DurationInt

object QueueTakeTimeoutCancelTest2 extends IOApp {

  override def run(args: List[String]): IO[ExitCode] =
    Queue
      .unbounded[IO, Int]
      .flatMap { queue =>
        val check: IO[Unit] =
          for {
            sink <- Deferred[IO, Int]

            // Goal:
            //  - Make sure that whenever an elem is dequeued, it is handled (at least via release).
            //  - Consumer must be cancelable, incl. when in `.take.timeout(...)`.
            consumer =
              IO.bracketFull(acquire = poll =>
                poll(queue.take.timeout(1.day)) // for the sake of this test, wait effectively forever
              )(
                use = elem => sink.complete(elem)
              )(
                release = (elem, _) => sink.complete(-elem).void
              )

            fib <- consumer.start

            // produce an element
            _ <- queue.offer(42)

            // and quickly cancel the consumer
            _ <- fib.cancel

            // inspect the outcome
            oc <- fib.join
            _ <- oc match {
              case Outcome.Succeeded(_) =>
                // succeeded to dequeue before cancel, check the value
                sink.tryGet.map(x => assert(x == Some(42)))
              case Outcome.Canceled() =>
                // Consumer was canceled.
                sink.tryGet.flatMap {
                  case Some(i) =>
                    // Consumer managed to dequeue the elem and it was not lost.
                    // sanity-check the value
                    IO.raiseWhen(i != 42 && i != -42) { new AssertionError(s"Unexpected value $i") }
                  case None =>
                    // Assuming consumer does not lose elements, the element must not have been dequeued.
                    // 1 second should be plenty of time to dequeue it (TimeoutException means element was lost)
                    queue.take.timeout(1.second)
                      .adaptError { case _: TimeoutException => new RuntimeException("Element was lost!") }
                      .void
                }
              case Outcome.Errored(e) =>
                IO.raiseError(new AssertionError("should never happen", e))
            }
          } yield ()

        check.replicateA_(10000)
      }
      .as(ExitCode.Success)
}

Output

java.lang.RuntimeException: Element was lost!
	at QueueTakeTimeoutCancelTest2$$anonfun$$nestedInanonfun$run$11$1.applyOrElse(QueueTakeTimeoutCancelTest2.scala:55)
	at QueueTakeTimeoutCancelTest2$$anonfun$$nestedInanonfun$run$11$1.applyOrElse(QueueTakeTimeoutCancelTest2.scala:55)
	at scala.PartialFunction$AndThen.applyOrElse(PartialFunction.scala:299)
	at timeout @ QueueTakeTimeoutCancelTest2$.$anonfun$run$3(QueueTakeTimeoutCancelTest2.scala:23)
	at main$ @ QueueTakeTimeoutCancelTest2$.main(QueueTakeTimeoutCancelTest2.scala:8)
	at main$ @ QueueTakeTimeoutCancelTest2$.main(QueueTakeTimeoutCancelTest2.scala:8)

Activity

  1. added theissue type on Mar 20, 2026
  2. durban commented on Mar 20, 2026

    @durban
    Contributor

    I don't think this is Queue specific. It's arguably not even timeout specific. I'm almost certain it is due to how racePair behaves.

    I'll still have to verify this, but I think what happens is this: RacePair have 3 fibers (the 2 forked ones, and the "current" fiber executing the RacePair). When there is only 1 .timeout(...), (i.e., the first example without the second timeout) the "current" fiber is never cancelled, and timeout handles correctly (since #4059) the case when the left side of the race already completed. However, when there are 2 .timeout(...)s (i.e., the first example as it is), then the inner RacePair itself is cancelled. And this cancellation don't/can't handle the case when the left side already completed (the relevant code is here: https://github.com/typelevel/cats-effect/blob/series/3.x/core/shared/src/main/scala/cats/effect/IOFiber.scala#L940).

    I think this issue is yet another manifestation of the fiber/cancellation model having some... issues; previous related tickets: #4489, notably #3553, #3555, and surely a few others I'm forgetting (e.g., the original "queue losing elements", which I can't find right now).

    I'm marking this as a bug, because it's really an unfortunate behavior. But I'm not sure how to fix this without changing how cancellation itself works (and even while #3553 looks good on the surface, it's exact semantics were not pinned down; that's partly why I've stopped working on #3917). Anyway, I didn't look at the general issue before specifically from the perspective of RacePair, so that's new (to me) :-) And maybe this specific case could be solved without the general one? (Like the single .timeout have been.)

  3. TomasMikula commented on Mar 20, 2026

    @TomasMikula
    Author

    Thanks for looking into this, @durban.

    This is surely very naive of me, but wouldn't calling the callback cb after joining the 2 canceled fibers if one of them succeeded cause the main, canceled fiber to in fact return successfully? Or has it by that time been irrevocably decided that the main fiber is canceled?

  4. durban commented on Mar 20, 2026

    @durban
    Contributor

    (No naïveté required, this is quite subtle stuff ;-) cancel in that code is a finalizer registered in the fiber executing racePair. That finalizer is only supposed to run, if that fiber is cancelled. If that fiber is already cancelled, calling cb mustn't have any effect. More precisely: if we're already cancelled and we're running finalizers, that means that the fiber was resumed (if it was suspended), so calling the callback will not cause it to be resumed again. In other words, this:

    by that time been irrevocably decided that the main fiber is canceled

    is exactly correct. (Assuming by "main fiber" you mean "the fiber executing the racePair".)

    Having said that, I admit I've not really experimented with this scenario so far (beyond confirming that I can reproduce it). My description above is based on my reading of the code (and previous experience debugging similar issues). In general I would encourage experimenting with this (if you feel inclined). (Also, thank you for the bug report, I think this shines a somewhat different light on issues we've seen before.)

  5. djspiewak commented on Jun 20, 2026

    @djspiewak
    Member

    I believe we can simplify this even further to the old chestnut:

    ioa.start.flatMap(inner => inner.join.onCancel(inner.cancel)).start.flatMap(outer => outer.cancel *> outer.join)

    If the outer fiber is canceled after inner completes but before outer completes, then inner can be canceled, tripping off the onCancel and voiding the result of outer (i.e. Outcome.Canceled()). However, inner has already completed, so the cancelation is spurious and the join was about to produce a result, but that result is already too late because we've pre-decided to set the inner outcome already.

    This race is fundamental, and you can contrive this same scenario without an inner fiber and instead some other async thing. I'm tempted to say (after thinking about it for like… 30 minutes) that this is kind of a classic "the ball falls between the two parties" type of thing. The cb completes inside of Cont, and at the moment it completes, we're canceled. We've written the semantics such that this deterministically resolves as prepareFiberForCancelation(null) followed by a resumption of the runloop (ultimately resulting in done(Outcome.Canceled())). The problem is that the value in the callback just kind of exists in this liminal state. It made it to state, so in theory the handoff is complete, but we ignored it and went ahead with cancelation.

    Put another way: inner has completed (which is why cb fired at all), and so it's no longer responsible for handling cancelation because it isn't being canceled, but outer doesn't realize that it's too late to cancel inner. It tries anyway, fails (doesn't realize this), and ignores the fact that it was actually handed a result.

    It's tempting to say that we can fix this by sanity checking state before we jump into cancelation land, but the problem is that we won't even attempt to cancel inner until we're in cancelation land (because it's set up by onCancel in the async implementation, and thus is a finalizer on the stack somewhere that we can't see). The original sins here are basically that cancelation (and the Cont continuation, equivalently) return Unit rather than Boolean, but also the fact that async's cancelation is non-primitive and Cont's mechanism is simply oblivious to it. This in turn makes me wonder if we can perhaps work around this very specific case (it wouldn't work more generally with folks registering their own onCancels) by making async primitive again, but that feels brittle.

    Another plausible idea is that we could add some logic here https://github.com/typelevel/cats-effect/blob/series/3.x/core/shared/src/main/scala/cats/effect/IOFiber.scala#L699 and here https://github.com/typelevel/cats-effect/blob/series/3.x/core/shared/src/main/scala/cats/effect/IOFiber.scala#L711 where we detect the fact that the callback was fired late. We would need to somehow save off the continuation stack prior to finalization though, otherwise there would be… nothing to do. And even in this case, we'd be talking about effectively forking a new zombie fiber while the old fiber merrily cancels itself.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    Type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions