Skip to content

Commit

Permalink
Preserve the relative order of delayed and non-delayed tasks in event…
Browse files Browse the repository at this point in the history
… loops

Fixed #4134
  • Loading branch information
dkhalanskyjb committed May 27, 2024
1 parent 2803a33 commit d2a21fc
Show file tree
Hide file tree
Showing 2 changed files with 37 additions and 15 deletions.
37 changes: 22 additions & 15 deletions kotlinx-coroutines-core/common/src/EventLoop.common.kt
Original file line number Diff line number Diff line change
Expand Up @@ -256,21 +256,7 @@ internal abstract class EventLoopImplBase: EventLoopImplPlatform(), Delay {
// unconfined events take priority
if (processUnconfinedEvent()) return 0
// queue all delayed tasks that are due to be executed
val delayed = _delayed.value
if (delayed != null && !delayed.isEmpty) {
val now = nanoTime()
while (true) {
// make sure that moving from delayed to queue removes from delayed only after it is added to queue
// to make sure that 'isEmpty' and `nextTime` that check both of them
// do not transiently report that both delayed and queue are empty during move
delayed.removeFirstIf {
if (it.timeToExecute(now)) {
enqueueImpl(it)
} else
false
} ?: break // quit loop when nothing more to remove or enqueueImpl returns false on "isComplete"
}
}
enqueueDelayedTasks()
// then process one event from queue
val task = dequeue()
if (task != null) {
Expand All @@ -283,6 +269,8 @@ internal abstract class EventLoopImplBase: EventLoopImplPlatform(), Delay {
final override fun dispatch(context: CoroutineContext, block: Runnable) = enqueue(block)

open fun enqueue(task: Runnable) {
// are there some delayed tasks that should execute before this one? If so, move them to the queue first.
enqueueDelayedTasks()
if (enqueueImpl(task)) {
// todo: we should unpark only when this delayed task became first in the queue
unpark()
Expand Down Expand Up @@ -336,6 +324,25 @@ internal abstract class EventLoopImplBase: EventLoopImplPlatform(), Delay {
}
}

/** Move all delayed tasks that are due to the main queue. */
private fun enqueueDelayedTasks() {
val delayed = _delayed.value
if (delayed != null && !delayed.isEmpty) {
val now = nanoTime()
while (true) {
// make sure that moving from delayed to queue removes from delayed only after it is added to queue
// to make sure that 'isEmpty' and `nextTime` that check both of them
// do not transiently report that both delayed and queue are empty during move
delayed.removeFirstIf {
if (it.timeToExecute(now)) {
enqueueImpl(it)
} else
false
} ?: break // quit loop when nothing more to remove or enqueueImpl returns false on "isComplete"
}
}
}

private fun closeQueue() {
assert { isCompleted }
_queue.loop { queue ->
Expand Down
15 changes: 15 additions & 0 deletions kotlinx-coroutines-core/jvm/test/EventLoopsTest.kt
Original file line number Diff line number Diff line change
Expand Up @@ -126,6 +126,21 @@ class EventLoopsTest : TestBase() {
finish(4)
}

/**
* Tests that, when delayed tasks are due on an event loop, they will execute earlier than the newly-scheduled
* non-delayed tasks.
*/
@Test
fun testPendingDelayedBeingDueEarlier() = runTest {
launch(start = CoroutineStart.UNDISPATCHED) {
delay(1)
expect(1)
}
Thread.sleep(100)
yield()
finish(2)
}

class EventSync {
private val waitingThread = atomic<Thread?>(null)
private val fired = atomic(false)
Expand Down

0 comments on commit d2a21fc

Please sign in to comment.