make asyncdispatch.poll completing all opterations that can be comple… (#6911)
introduce asyncdispatch.drain that completes all operations that can be completed immediately; implements #6523
This commit is contained in:
parent
3de81af44d
commit
85ac3130aa
2 changed files with 35 additions and 14 deletions
|
|
@ -168,18 +168,20 @@ type
|
||||||
timers*: HeapQueue[tuple[finishAt: float, fut: Future[void]]]
|
timers*: HeapQueue[tuple[finishAt: float, fut: Future[void]]]
|
||||||
callbacks*: Deque[proc ()]
|
callbacks*: Deque[proc ()]
|
||||||
|
|
||||||
proc processTimers(p: PDispatcherBase) {.inline.} =
|
proc processTimers(p: PDispatcherBase; didSomeWork: var bool) {.inline.} =
|
||||||
#Process just part if timers at a step
|
#Process just part if timers at a step
|
||||||
var count = p.timers.len
|
var count = p.timers.len
|
||||||
let t = epochTime()
|
let t = epochTime()
|
||||||
while count > 0 and t >= p.timers[0].finishAt:
|
while count > 0 and t >= p.timers[0].finishAt:
|
||||||
p.timers.pop().fut.complete()
|
p.timers.pop().fut.complete()
|
||||||
dec count
|
dec count
|
||||||
|
didSomeWork = true
|
||||||
|
|
||||||
proc processPendingCallbacks(p: PDispatcherBase) =
|
proc processPendingCallbacks(p: PDispatcherBase; didSomeWork: var bool) =
|
||||||
while p.callbacks.len > 0:
|
while p.callbacks.len > 0:
|
||||||
var cb = p.callbacks.popFirst()
|
var cb = p.callbacks.popFirst()
|
||||||
cb()
|
cb()
|
||||||
|
didSomeWork = true
|
||||||
|
|
||||||
proc adjustedTimeout(p: PDispatcherBase, timeout: int): int {.inline.} =
|
proc adjustedTimeout(p: PDispatcherBase, timeout: int): int {.inline.} =
|
||||||
# If dispatcher has active timers this proc returns the timeout
|
# If dispatcher has active timers this proc returns the timeout
|
||||||
|
|
@ -284,14 +286,13 @@ when defined(windows) or defined(nimdoc):
|
||||||
let p = getGlobalDispatcher()
|
let p = getGlobalDispatcher()
|
||||||
p.handles.len != 0 or p.timers.len != 0 or p.callbacks.len != 0
|
p.handles.len != 0 or p.timers.len != 0 or p.callbacks.len != 0
|
||||||
|
|
||||||
proc poll*(timeout = 500) =
|
proc runOnce(timeout = 500): bool =
|
||||||
## Waits for completion events and processes them. Raises ``ValueError``
|
|
||||||
## if there are no pending operations.
|
|
||||||
let p = getGlobalDispatcher()
|
let p = getGlobalDispatcher()
|
||||||
if p.handles.len == 0 and p.timers.len == 0 and p.callbacks.len == 0:
|
if p.handles.len == 0 and p.timers.len == 0 and p.callbacks.len == 0:
|
||||||
raise newException(ValueError,
|
raise newException(ValueError,
|
||||||
"No handles or timers registered in dispatcher.")
|
"No handles or timers registered in dispatcher.")
|
||||||
|
|
||||||
|
result = false
|
||||||
if p.handles.len != 0:
|
if p.handles.len != 0:
|
||||||
let at = p.adjustedTimeout(timeout)
|
let at = p.adjustedTimeout(timeout)
|
||||||
var llTimeout =
|
var llTimeout =
|
||||||
|
|
@ -304,6 +305,7 @@ when defined(windows) or defined(nimdoc):
|
||||||
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
|
||||||
|
|
||||||
# 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
|
||||||
|
|
@ -333,13 +335,14 @@ when defined(windows) or defined(nimdoc):
|
||||||
else:
|
else:
|
||||||
if errCode.int32 == WAIT_TIMEOUT:
|
if errCode.int32 == WAIT_TIMEOUT:
|
||||||
# Timed out
|
# Timed out
|
||||||
discard
|
result = false
|
||||||
else: raiseOSError(errCode)
|
else: raiseOSError(errCode)
|
||||||
|
|
||||||
# Timer processing.
|
# Timer processing.
|
||||||
processTimers(p)
|
processTimers(p, result)
|
||||||
# Callback queue processing
|
# Callback queue processing
|
||||||
processPendingCallbacks(p)
|
processPendingCallbacks(p, result)
|
||||||
|
|
||||||
|
|
||||||
var acceptEx: WSAPROC_ACCEPTEX
|
var acceptEx: WSAPROC_ACCEPTEX
|
||||||
var connectEx: WSAPROC_CONNECTEX
|
var connectEx: WSAPROC_CONNECTEX
|
||||||
|
|
@ -1202,7 +1205,7 @@ else:
|
||||||
# descriptor was unregistered in callback via `unregister()`.
|
# descriptor was unregistered in callback via `unregister()`.
|
||||||
discard
|
discard
|
||||||
|
|
||||||
proc poll*(timeout = 500) =
|
proc runOnce(timeout = 500): bool =
|
||||||
let p = getGlobalDispatcher()
|
let p = getGlobalDispatcher()
|
||||||
when ioselSupportedPlatform:
|
when ioselSupportedPlatform:
|
||||||
let customSet = {Event.Timer, Event.Signal, Event.Process,
|
let customSet = {Event.Timer, Event.Signal, Event.Process,
|
||||||
|
|
@ -1212,6 +1215,7 @@ else:
|
||||||
raise newException(ValueError,
|
raise newException(ValueError,
|
||||||
"No handles or timers registered in dispatcher.")
|
"No handles or timers registered in dispatcher.")
|
||||||
|
|
||||||
|
result = false
|
||||||
if not p.selector.isEmpty():
|
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)
|
||||||
|
|
@ -1224,20 +1228,24 @@ else:
|
||||||
|
|
||||||
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
|
||||||
|
|
||||||
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
|
||||||
|
|
||||||
if Event.User in events or events == {Event.Error}:
|
if Event.User in events or events == {Event.Error}:
|
||||||
processBasicCallbacks(fd, readList)
|
processBasicCallbacks(fd, readList)
|
||||||
custom = true
|
custom = true
|
||||||
if rLength == 0:
|
if rLength == 0:
|
||||||
p.selector.unregister(fd)
|
p.selector.unregister(fd)
|
||||||
|
result = true
|
||||||
|
|
||||||
when ioselSupportedPlatform:
|
when ioselSupportedPlatform:
|
||||||
if (customSet * events) != {}:
|
if (customSet * events) != {}:
|
||||||
custom = true
|
custom = true
|
||||||
processCustomCallbacks(fd)
|
processCustomCallbacks(fd)
|
||||||
|
result = true
|
||||||
|
|
||||||
# because state `data` can be modified in callback we need to update
|
# because state `data` can be modified in callback we need to update
|
||||||
# descriptor events with currently registered callbacks.
|
# descriptor events with currently registered callbacks.
|
||||||
|
|
@ -1249,9 +1257,9 @@ else:
|
||||||
p.selector.updateHandle(SocketHandle(fd), newEvents)
|
p.selector.updateHandle(SocketHandle(fd), newEvents)
|
||||||
|
|
||||||
# Timer processing.
|
# Timer processing.
|
||||||
processTimers(p)
|
processTimers(p, result)
|
||||||
# Callback queue processing
|
# Callback queue processing
|
||||||
processPendingCallbacks(p)
|
processPendingCallbacks(p, result)
|
||||||
|
|
||||||
proc recv*(socket: AsyncFD, size: int,
|
proc recv*(socket: AsyncFD, size: int,
|
||||||
flags = {SocketFlag.SafeDisconn}): Future[string] =
|
flags = {SocketFlag.SafeDisconn}): Future[string] =
|
||||||
|
|
@ -1474,6 +1482,19 @@ else:
|
||||||
data.readList.add(cb)
|
data.readList.add(cb)
|
||||||
p.selector.registerEvent(SelectEvent(ev), data)
|
p.selector.registerEvent(SelectEvent(ev), data)
|
||||||
|
|
||||||
|
proc drain*(timeout = 500) =
|
||||||
|
## Waits for completion events and processes them. Raises ``ValueError``
|
||||||
|
## if there are no pending operations. In contrast to ``poll`` this
|
||||||
|
## processes as many events as are available.
|
||||||
|
if runOnce(timeout):
|
||||||
|
while runOnce(0): discard
|
||||||
|
|
||||||
|
proc poll*(timeout = 500) =
|
||||||
|
## Waits for completion events and processes them. Raises ``ValueError``
|
||||||
|
## if there are no pending operations. This runs the underlying OS
|
||||||
|
## `epoll`:idx: or `kqueue`:idx: primitive only once.
|
||||||
|
discard runOnce()
|
||||||
|
|
||||||
# Common procedures between current and upcoming asyncdispatch
|
# Common procedures between current and upcoming asyncdispatch
|
||||||
include includes.asynccommon
|
include includes.asynccommon
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -579,9 +579,9 @@ else:
|
||||||
var event = newSelectEvent()
|
var event = newSelectEvent()
|
||||||
selector.registerEvent(event, 1)
|
selector.registerEvent(event, 1)
|
||||||
discard selector.select(0)
|
discard selector.select(0)
|
||||||
event.setEvent()
|
event.trigger()
|
||||||
var rc1 = selector.select(0)
|
var rc1 = selector.select(0)
|
||||||
event.setEvent()
|
event.trigger()
|
||||||
var rc2 = selector.select(0)
|
var rc2 = selector.select(0)
|
||||||
var rc3 = selector.select(0)
|
var rc3 = selector.select(0)
|
||||||
assert(len(rc1) == 1 and len(rc2) == 1 and len(rc3) == 0)
|
assert(len(rc1) == 1 and len(rc2) == 1 and len(rc3) == 0)
|
||||||
|
|
@ -611,7 +611,7 @@ else:
|
||||||
var event = newSelectEvent()
|
var event = newSelectEvent()
|
||||||
for i in 0..high(thr):
|
for i in 0..high(thr):
|
||||||
createThread(thr[i], event_wait_thread, event)
|
createThread(thr[i], event_wait_thread, event)
|
||||||
event.setEvent()
|
event.trigger()
|
||||||
joinThreads(thr)
|
joinThreads(thr)
|
||||||
assert(counter == 1)
|
assert(counter == 1)
|
||||||
result = true
|
result = true
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue