Fixes #4262.
This commit is contained in:
parent
04c1caf025
commit
5bf16439e1
3 changed files with 94 additions and 68 deletions
|
|
@ -500,48 +500,49 @@ when defined(windows) or defined(nimdoc):
|
||||||
raise newException(ValueError,
|
raise newException(ValueError,
|
||||||
"No handles or timers registered in dispatcher.")
|
"No handles or timers registered in dispatcher.")
|
||||||
|
|
||||||
let at = p.adjustedTimeout(timeout)
|
if p.handles.len != 0:
|
||||||
var llTimeout =
|
let at = p.adjustedTimeout(timeout)
|
||||||
if at == -1: winlean.INFINITE
|
var llTimeout =
|
||||||
else: at.int32
|
if at == -1: winlean.INFINITE
|
||||||
|
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
|
||||||
|
|
||||||
# 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, errCode)
|
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:
|
if customOverlapped.data.cell.data != nil:
|
||||||
system.dispose(customOverlapped.data.cell)
|
system.dispose(customOverlapped.data.cell)
|
||||||
|
|
||||||
GC_unref(customOverlapped)
|
GC_unref(customOverlapped)
|
||||||
else:
|
else:
|
||||||
if errCode.int32 == WAIT_TIMEOUT:
|
let errCode = osLastError()
|
||||||
# Timed out
|
if customOverlapped != nil:
|
||||||
discard
|
assert customOverlapped.data.fd == lpCompletionKey.AsyncFD
|
||||||
else: raiseOSError(errCode)
|
customOverlapped.data.cb(customOverlapped.data.fd,
|
||||||
|
lpNumberOfBytesTransferred, errCode)
|
||||||
|
if customOverlapped.data.cell.data != nil:
|
||||||
|
system.dispose(customOverlapped.data.cell)
|
||||||
|
GC_unref(customOverlapped)
|
||||||
|
else:
|
||||||
|
if errCode.int32 == WAIT_TIMEOUT:
|
||||||
|
# Timed out
|
||||||
|
discard
|
||||||
|
else: raiseOSError(errCode)
|
||||||
|
|
||||||
# Timer processing.
|
# Timer processing.
|
||||||
processTimers(p)
|
processTimers(p)
|
||||||
|
|
@ -1283,43 +1284,45 @@ else:
|
||||||
|
|
||||||
proc poll*(timeout = 500) =
|
proc poll*(timeout = 500) =
|
||||||
let p = getGlobalDispatcher()
|
let p = getGlobalDispatcher()
|
||||||
for info in p.selector.select(p.adjustedTimeout(timeout)):
|
|
||||||
let data = PData(info.key.data)
|
|
||||||
assert data.fd == info.key.fd.AsyncFD
|
|
||||||
#echo("In poll ", data.fd.cint)
|
|
||||||
# There may be EvError here, but we handle them in callbacks,
|
|
||||||
# so that exceptions can be raised from `send(...)` and
|
|
||||||
# `recv(...)` routines.
|
|
||||||
|
|
||||||
if EvRead in info.events:
|
if p.selector.len > 0:
|
||||||
# Callback may add items to ``data.readCBs`` which causes issues if
|
for info in p.selector.select(p.adjustedTimeout(timeout)):
|
||||||
# we are iterating over ``data.readCBs`` at the same time. We therefore
|
let data = PData(info.key.data)
|
||||||
# make a copy to iterate over.
|
assert data.fd == info.key.fd.AsyncFD
|
||||||
let currentCBs = data.readCBs
|
#echo("In poll ", data.fd.cint)
|
||||||
data.readCBs = @[]
|
# There may be EvError here, but we handle them in callbacks,
|
||||||
for cb in currentCBs:
|
# so that exceptions can be raised from `send(...)` and
|
||||||
if not cb(data.fd):
|
# `recv(...)` routines.
|
||||||
# Callback wants to be called again.
|
|
||||||
data.readCBs.add(cb)
|
|
||||||
|
|
||||||
if EvWrite in info.events:
|
if EvRead in info.events:
|
||||||
let currentCBs = data.writeCBs
|
# Callback may add items to ``data.readCBs`` which causes issues if
|
||||||
data.writeCBs = @[]
|
# we are iterating over ``data.readCBs`` at the same time. We therefore
|
||||||
for cb in currentCBs:
|
# make a copy to iterate over.
|
||||||
if not cb(data.fd):
|
let currentCBs = data.readCBs
|
||||||
# Callback wants to be called again.
|
data.readCBs = @[]
|
||||||
data.writeCBs.add(cb)
|
for cb in currentCBs:
|
||||||
|
if not cb(data.fd):
|
||||||
|
# Callback wants to be called again.
|
||||||
|
data.readCBs.add(cb)
|
||||||
|
|
||||||
if info.key in p.selector:
|
if EvWrite in info.events:
|
||||||
var newEvents: set[Event]
|
let currentCBs = data.writeCBs
|
||||||
if data.readCBs.len != 0: newEvents = {EvRead}
|
data.writeCBs = @[]
|
||||||
if data.writeCBs.len != 0: newEvents = newEvents + {EvWrite}
|
for cb in currentCBs:
|
||||||
if newEvents != info.key.events:
|
if not cb(data.fd):
|
||||||
update(data.fd, newEvents)
|
# Callback wants to be called again.
|
||||||
else:
|
data.writeCBs.add(cb)
|
||||||
# FD no longer a part of the selector. Likely been closed
|
|
||||||
# (e.g. socket disconnected).
|
if info.key in p.selector:
|
||||||
discard
|
var newEvents: set[Event]
|
||||||
|
if data.readCBs.len != 0: newEvents = {EvRead}
|
||||||
|
if data.writeCBs.len != 0: newEvents = newEvents + {EvWrite}
|
||||||
|
if newEvents != info.key.events:
|
||||||
|
update(data.fd, newEvents)
|
||||||
|
else:
|
||||||
|
# FD no longer a part of the selector. Likely been closed
|
||||||
|
# (e.g. socket disconnected).
|
||||||
|
discard
|
||||||
|
|
||||||
# Timer processing.
|
# Timer processing.
|
||||||
processTimers(p)
|
processTimers(p)
|
||||||
|
|
|
||||||
|
|
@ -375,6 +375,10 @@ proc contains*(s: Selector, key: SelectorKey): bool =
|
||||||
when not defined(nimdoc):
|
when not defined(nimdoc):
|
||||||
return key.fd in s and s.fds[key.fd] == key
|
return key.fd in s and s.fds[key.fd] == key
|
||||||
|
|
||||||
|
proc len*(s: Selector): int =
|
||||||
|
## Retrieves the number of registered file descriptors in this Selector.
|
||||||
|
return s.fds.len
|
||||||
|
|
||||||
{.deprecated: [TEvent: Event, PSelectorKey: SelectorKey,
|
{.deprecated: [TEvent: Event, PSelectorKey: SelectorKey,
|
||||||
TReadyInfo: ReadyInfo, PSelector: Selector].}
|
TReadyInfo: ReadyInfo, PSelector: Selector].}
|
||||||
|
|
||||||
|
|
|
||||||
19
tests/async/tpolltimeouts.nim
Normal file
19
tests/async/tpolltimeouts.nim
Normal file
|
|
@ -0,0 +1,19 @@
|
||||||
|
discard """
|
||||||
|
output: "true"
|
||||||
|
"""
|
||||||
|
# Issue https://github.com/nim-lang/Nim/issues/4262
|
||||||
|
import asyncdispatch, times
|
||||||
|
|
||||||
|
proc foo(): Future[int] {.async.} =
|
||||||
|
return 1
|
||||||
|
|
||||||
|
proc bar(): Future[int] {.async.} =
|
||||||
|
return await foo()
|
||||||
|
|
||||||
|
let start = epochTime()
|
||||||
|
let barFut = bar()
|
||||||
|
|
||||||
|
while not barFut.finished:
|
||||||
|
poll(2000)
|
||||||
|
|
||||||
|
echo(epochTime() - start < 1.0)
|
||||||
Loading…
Add table
Add a link
Reference in a new issue