diff --git a/lib/pure/asyncdispatch.nim b/lib/pure/asyncdispatch.nim index 70d94b023e..3c0a2257b0 100644 --- a/lib/pure/asyncdispatch.nim +++ b/lib/pure/asyncdispatch.nim @@ -269,6 +269,23 @@ proc processPendingCallbacks(p: PDispatcherBase; didSomeWork: var bool) = cb() didSomeWork = true +proc processTimersBeforePoll( + p: PDispatcherBase, didSomeWork: var bool +): Option[int] {.inline.} = + # Do not let an expired timeout overtake completion callbacks which are + # already pending. `adjustTimeout` makes the I/O poll non-blocking when the + # callback queue is non-empty. + if p.callbacks.len == 0: + result = processTimers(p, didSomeWork) + +proc processCallbacksAndTimers(p: PDispatcherBase; didSomeWork: var bool) = + # A completed operation can take multiple queued callbacks to propagate + # through its public future. Process the whole chain before expired timers. + processPendingCallbacks(p, didSomeWork) + discard processTimers(p, didSomeWork) + # Timer futures must still propagate within this dispatcher iteration. + processPendingCallbacks(p, didSomeWork) + proc adjustTimeout( p: PDispatcherBase, pollTimeout: int, nextTimer: Option[int] ): int {.inline.} = @@ -399,7 +416,7 @@ when defined(windows) or defined(nimdoc): "No handles or timers registered in dispatcher.") result = false - let nextTimer = processTimers(p, result) + let nextTimer = processTimersBeforePoll(p, result) let at = adjustTimeout(p, timeout, nextTimer) var llTimeout = if at == -1: winlean.INFINITE @@ -450,10 +467,7 @@ when defined(windows) or defined(nimdoc): result = false else: raiseOSError(errCode) - # Timer processing. - discard processTimers(p, result) - # Callback queue processing - processPendingCallbacks(p, result) + processCallbacksAndTimers(p, result) var acceptEx: WSAPROC_ACCEPTEX @@ -1404,7 +1418,7 @@ else: result = false var keys: array[64, ReadyKey] - let nextTimer = processTimers(p, result) + let nextTimer = processTimersBeforePoll(p, result) var count = p.selector.selectInto(adjustTimeout(p, timeout, nextTimer), keys) for i in 0.. 0: incl(newEvents, Event.Write) p.selector.updateHandle(SocketHandle(fd), newEvents) - # Timer processing. - discard processTimers(p, result) - # Callback queue processing - processPendingCallbacks(p, result) + processCallbacksAndTimers(p, result) proc recv*(socket: AsyncFD, size: int, flags = {SocketFlag.SafeDisconn}): owned(Future[string]) = diff --git a/tests/async/tasyncclosestall.nim b/tests/async/tasyncclosestall.nim index 05348587c1..ea4865fc79 100644 --- a/tests/async/tasyncclosestall.nim +++ b/tests/async/tasyncclosestall.nim @@ -4,6 +4,7 @@ discard """ exitcode: 0 """ import asyncdispatch, asyncnet +import std/strutils when defined(windows): from winlean import ERROR_NETNAME_DELETED @@ -14,6 +15,7 @@ else: # even when the socket is closed. const timeout = 2000 + messagePaddingSize = 64 * 1024 var port = Port(0) var sent = 0 @@ -31,10 +33,12 @@ proc isExpectedDisconnectionError(errCode: int32): bool = errCode == EBADF or errCode == ECONNRESET or errCode == EPIPE proc keepSendingTo(c: AsyncSocket) {.async.} = + let messagePadding = repeat('x', messagePaddingSize) while true: - # This write will eventually get stuck because the client is not reading - # its messages. - let sendFut = c.send("Foobar" & $sent & "\n", flags = {}) + # Larger writes reach socket backpressure quickly even on slow CI machines. + # This write will eventually get stuck because the client is not reading. + # Keep the padding after the newline so recvLine does not drain it. + let sendFut = c.send("Foobar" & $sent & "\n" & messagePadding, flags = {}) var sendTimedOut = false try: # On some platforms (notably macOS ARM64), the kernel may return diff --git a/tests/async/tasyncdispatchordering.nim b/tests/async/tasyncdispatchordering.nim new file mode 100644 index 0000000000..5b226bad8a --- /dev/null +++ b/tests/async/tasyncdispatchordering.nim @@ -0,0 +1,32 @@ +discard """ + action: run +""" + +import asyncdispatch, os + +proc wrap(fut: Future[void]): Future[void] = + result = newFuture[void]("wrap") + let retFuture = result + fut.addCallback proc () = + if fut.failed: + retFuture.fail(fut.error) + else: + retFuture.complete() + +block: + let root = newFuture[void]("root") + let wrapped = wrap(wrap(wrap(root))) + let completedBeforeDeadline = withTimeout(wrapped, 20) + + # Completion has happened at the bottom of the future chain, but its + # callbacks cannot propagate until control reaches the dispatcher. + root.complete() + sleep(40) + + doAssert waitFor(completedBeforeDeadline) + +block: + var callbackRan = false + sleepAsync(0).addCallback proc () = callbackRan = true + poll(0) + doAssert callbackRan