Don't skip poll() when no handles are present. (#8727)
Fixes #7886. Fixes #7758. Fixes #6929. Fixes #3909. Replaces #8209.
This commit is contained in:
parent
55a8649749
commit
7532b37405
2 changed files with 84 additions and 69 deletions
|
|
@ -299,50 +299,49 @@ when defined(windows) or defined(nimdoc):
|
||||||
"No handles or timers registered in dispatcher.")
|
"No handles or timers registered in dispatcher.")
|
||||||
|
|
||||||
result = false
|
result = false
|
||||||
if p.handles.len != 0:
|
let at = p.adjustedTimeout(timeout)
|
||||||
let at = p.adjustedTimeout(timeout)
|
var llTimeout =
|
||||||
var llTimeout =
|
if at == -1: winlean.INFINITE
|
||||||
if at == -1: winlean.INFINITE
|
else: at.int32
|
||||||
else: at.int32
|
|
||||||
|
|
||||||
var lpNumberOfBytesTransferred: Dword
|
var lpNumberOfBytesTransferred: Dword
|
||||||
var lpCompletionKey: ULONG_PTR
|
var lpCompletionKey: ULONG_PTR
|
||||||
var customOverlapped: PCustomOverlapped
|
var customOverlapped: PCustomOverlapped
|
||||||
let res = getQueuedCompletionStatus(p.ioPort,
|
let res = getQueuedCompletionStatus(p.ioPort,
|
||||||
addr lpNumberOfBytesTransferred, addr lpCompletionKey,
|
addr lpNumberOfBytesTransferred, addr lpCompletionKey,
|
||||||
cast[ptr POVERLAPPED](addr customOverlapped), llTimeout).bool
|
cast[ptr POVERLAPPED](addr customOverlapped), llTimeout).bool
|
||||||
result = true
|
result = true
|
||||||
|
|
||||||
# http://stackoverflow.com/a/12277264/492186
|
# http://stackoverflow.com/a/12277264/492186
|
||||||
# TODO: http://www.serverframework.com/handling-multiple-pending-socket-read-and-write-operations.html
|
# TODO: http://www.serverframework.com/handling-multiple-pending-socket-read-and-write-operations.html
|
||||||
if res:
|
if res:
|
||||||
# This is useful for ensuring the reliability of the overlapped struct.
|
# This is useful for ensuring the reliability of the overlapped struct.
|
||||||
|
assert customOverlapped.data.fd == lpCompletionKey.AsyncFD
|
||||||
|
|
||||||
|
customOverlapped.data.cb(customOverlapped.data.fd,
|
||||||
|
lpNumberOfBytesTransferred, OSErrorCode(-1))
|
||||||
|
|
||||||
|
# If cell.data != nil, then system.protect(rawEnv(cb)) was called,
|
||||||
|
# so we need to dispose our `cb` environment, because it is not needed
|
||||||
|
# anymore.
|
||||||
|
if customOverlapped.data.cell.data != nil:
|
||||||
|
system.dispose(customOverlapped.data.cell)
|
||||||
|
|
||||||
|
GC_unref(customOverlapped)
|
||||||
|
else:
|
||||||
|
let errCode = osLastError()
|
||||||
|
if customOverlapped != nil:
|
||||||
assert customOverlapped.data.fd == lpCompletionKey.AsyncFD
|
assert customOverlapped.data.fd == lpCompletionKey.AsyncFD
|
||||||
|
|
||||||
customOverlapped.data.cb(customOverlapped.data.fd,
|
customOverlapped.data.cb(customOverlapped.data.fd,
|
||||||
lpNumberOfBytesTransferred, OSErrorCode(-1))
|
lpNumberOfBytesTransferred, errCode)
|
||||||
|
|
||||||
# If cell.data != nil, then system.protect(rawEnv(cb)) was called,
|
|
||||||
# so we need to dispose our `cb` environment, because it is not needed
|
|
||||||
# anymore.
|
|
||||||
if customOverlapped.data.cell.data != nil:
|
if customOverlapped.data.cell.data != nil:
|
||||||
system.dispose(customOverlapped.data.cell)
|
system.dispose(customOverlapped.data.cell)
|
||||||
|
|
||||||
GC_unref(customOverlapped)
|
GC_unref(customOverlapped)
|
||||||
else:
|
else:
|
||||||
let errCode = osLastError()
|
if errCode.int32 == WAIT_TIMEOUT:
|
||||||
if customOverlapped != nil:
|
# Timed out
|
||||||
assert customOverlapped.data.fd == lpCompletionKey.AsyncFD
|
result = false
|
||||||
customOverlapped.data.cb(customOverlapped.data.fd,
|
else: raiseOSError(errCode)
|
||||||
lpNumberOfBytesTransferred, errCode)
|
|
||||||
if customOverlapped.data.cell.data != nil:
|
|
||||||
system.dispose(customOverlapped.data.cell)
|
|
||||||
GC_unref(customOverlapped)
|
|
||||||
else:
|
|
||||||
if errCode.int32 == WAIT_TIMEOUT:
|
|
||||||
# Timed out
|
|
||||||
result = false
|
|
||||||
else: raiseOSError(errCode)
|
|
||||||
|
|
||||||
# Timer processing.
|
# Timer processing.
|
||||||
processTimers(p, result)
|
processTimers(p, result)
|
||||||
|
|
@ -1231,45 +1230,44 @@ else:
|
||||||
"No handles or timers registered in dispatcher.")
|
"No handles or timers registered in dispatcher.")
|
||||||
|
|
||||||
result = false
|
result = false
|
||||||
if not p.selector.isEmpty():
|
var keys: array[64, ReadyKey]
|
||||||
var keys: array[64, ReadyKey]
|
var count = p.selector.selectInto(p.adjustedTimeout(timeout), keys)
|
||||||
var count = p.selector.selectInto(p.adjustedTimeout(timeout), keys)
|
for i in 0..<count:
|
||||||
for i in 0..<count:
|
var custom = false
|
||||||
var custom = false
|
let fd = keys[i].fd
|
||||||
let fd = keys[i].fd
|
let events = keys[i].events
|
||||||
let events = keys[i].events
|
var rLength = 0 # len(data.readList) after callback
|
||||||
var rLength = 0 # len(data.readList) after callback
|
var wLength = 0 # len(data.writeList) after callback
|
||||||
var wLength = 0 # len(data.writeList) after callback
|
|
||||||
|
|
||||||
if Event.Read in events or events == {Event.Error}:
|
if Event.Read in events or events == {Event.Error}:
|
||||||
processBasicCallbacks(fd, readList)
|
processBasicCallbacks(fd, readList)
|
||||||
result = true
|
result = true
|
||||||
|
|
||||||
if Event.Write in events or events == {Event.Error}:
|
if Event.Write in events or events == {Event.Error}:
|
||||||
processBasicCallbacks(fd, writeList)
|
processBasicCallbacks(fd, writeList)
|
||||||
result = true
|
result = true
|
||||||
|
|
||||||
if Event.User in events:
|
if Event.User in events:
|
||||||
processBasicCallbacks(fd, readList)
|
processBasicCallbacks(fd, readList)
|
||||||
|
custom = true
|
||||||
|
if rLength == 0:
|
||||||
|
p.selector.unregister(fd)
|
||||||
|
result = true
|
||||||
|
|
||||||
|
when ioselSupportedPlatform:
|
||||||
|
if (customSet * events) != {}:
|
||||||
custom = true
|
custom = true
|
||||||
if rLength == 0:
|
processCustomCallbacks(fd)
|
||||||
p.selector.unregister(fd)
|
|
||||||
result = true
|
result = true
|
||||||
|
|
||||||
when ioselSupportedPlatform:
|
# because state `data` can be modified in callback we need to update
|
||||||
if (customSet * events) != {}:
|
# descriptor events with currently registered callbacks.
|
||||||
custom = true
|
if not custom:
|
||||||
processCustomCallbacks(fd)
|
var newEvents: set[Event] = {}
|
||||||
result = true
|
if rLength != -1 and wLength != -1:
|
||||||
|
if rLength > 0: incl(newEvents, Event.Read)
|
||||||
# because state `data` can be modified in callback we need to update
|
if wLength > 0: incl(newEvents, Event.Write)
|
||||||
# descriptor events with currently registered callbacks.
|
p.selector.updateHandle(SocketHandle(fd), newEvents)
|
||||||
if not custom:
|
|
||||||
var newEvents: set[Event] = {}
|
|
||||||
if rLength != -1 and wLength != -1:
|
|
||||||
if rLength > 0: incl(newEvents, Event.Read)
|
|
||||||
if wLength > 0: incl(newEvents, Event.Write)
|
|
||||||
p.selector.updateHandle(SocketHandle(fd), newEvents)
|
|
||||||
|
|
||||||
# Timer processing.
|
# Timer processing.
|
||||||
processTimers(p, result)
|
processTimers(p, result)
|
||||||
|
|
|
||||||
17
tests/async/t7758.nim
Normal file
17
tests/async/t7758.nim
Normal file
|
|
@ -0,0 +1,17 @@
|
||||||
|
discard """
|
||||||
|
file: "t7758.nim"
|
||||||
|
exitcode: 0
|
||||||
|
"""
|
||||||
|
import asyncdispatch
|
||||||
|
|
||||||
|
proc task() {.async.} =
|
||||||
|
await sleepAsync(1000)
|
||||||
|
|
||||||
|
when isMainModule:
|
||||||
|
var counter = 0
|
||||||
|
var f = task()
|
||||||
|
while not f.finished:
|
||||||
|
inc(counter)
|
||||||
|
poll()
|
||||||
|
|
||||||
|
doAssert counter == 2
|
||||||
Loading…
Add table
Add a link
Reference in a new issue