async: minor refactorings (#15354)
This commit is contained in:
parent
e56d50d747
commit
2671efab78
5 changed files with 43 additions and 62 deletions
|
|
@ -290,7 +290,7 @@ when defined(windows) or defined(nimdoc):
|
||||||
|
|
||||||
var gDisp{.threadvar.}: owned PDispatcher ## Global dispatcher
|
var gDisp{.threadvar.}: owned PDispatcher ## Global dispatcher
|
||||||
|
|
||||||
proc setGlobalDispatcher*(disp: owned PDispatcher) =
|
proc setGlobalDispatcher*(disp: sink PDispatcher) =
|
||||||
if not gDisp.isNil:
|
if not gDisp.isNil:
|
||||||
assert gDisp.callbacks.len == 0
|
assert gDisp.callbacks.len == 0
|
||||||
gDisp = disp
|
gDisp = disp
|
||||||
|
|
@ -1217,10 +1217,12 @@ else:
|
||||||
withData(selector, fd.int, fdData):
|
withData(selector, fd.int, fdData):
|
||||||
case event
|
case event
|
||||||
of Event.Read:
|
of Event.Read:
|
||||||
shallowCopy(curList, fdData.readList)
|
#shallowCopy(curList, fdData.readList)
|
||||||
|
curList = move fdData.readList
|
||||||
fdData.readList = newSeqOfCap[Callback](InitCallbackListSize)
|
fdData.readList = newSeqOfCap[Callback](InitCallbackListSize)
|
||||||
of Event.Write:
|
of Event.Write:
|
||||||
shallowCopy(curList, fdData.writeList)
|
#shallowCopy(curList, fdData.writeList)
|
||||||
|
curList = move fdData.writeList
|
||||||
fdData.writeList = newSeqOfCap[Callback](InitCallbackListSize)
|
fdData.writeList = newSeqOfCap[Callback](InitCallbackListSize)
|
||||||
else:
|
else:
|
||||||
assert false, "Cannot process callbacks for " & $event
|
assert false, "Cannot process callbacks for " & $event
|
||||||
|
|
@ -1232,8 +1234,7 @@ else:
|
||||||
for cb in curList:
|
for cb in curList:
|
||||||
if eventsExtinguished:
|
if eventsExtinguished:
|
||||||
newList.add(cb)
|
newList.add(cb)
|
||||||
continue
|
elif not cb(fd):
|
||||||
if not cb(fd):
|
|
||||||
# Callback wants to be called again.
|
# Callback wants to be called again.
|
||||||
newList.add(cb)
|
newList.add(cb)
|
||||||
# This callback has returned with EAGAIN, so we don't need to
|
# This callback has returned with EAGAIN, so we don't need to
|
||||||
|
|
@ -1259,15 +1260,15 @@ else:
|
||||||
result.readCbListCount = -1
|
result.readCbListCount = -1
|
||||||
result.writeCbListCount = -1
|
result.writeCbListCount = -1
|
||||||
|
|
||||||
template processCustomCallbacks(ident: untyped) =
|
proc processCustomCallbacks(p: PDispatcher; fd: AsyncFD) =
|
||||||
# Process pending custom event callbacks. Custom events are
|
# Process pending custom event callbacks. Custom events are
|
||||||
# {Event.Timer, Event.Signal, Event.Process, Event.Vnode}.
|
# {Event.Timer, Event.Signal, Event.Process, Event.Vnode}.
|
||||||
# There can be only one callback registered with one descriptor,
|
# There can be only one callback registered with one descriptor,
|
||||||
# so there is no need to iterate over list.
|
# so there is no need to iterate over list.
|
||||||
var curList: seq[Callback]
|
var curList: seq[Callback]
|
||||||
|
|
||||||
withData(p.selector, ident.int, adata) do:
|
withData(p.selector, fd.int, adata) do:
|
||||||
shallowCopy(curList, adata.readList)
|
curList = move adata.readList
|
||||||
adata.readList = newSeqOfCap[Callback](InitCallbackListSize)
|
adata.readList = newSeqOfCap[Callback](InitCallbackListSize)
|
||||||
|
|
||||||
let newLength = len(curList)
|
let newLength = len(curList)
|
||||||
|
|
@ -1277,7 +1278,7 @@ else:
|
||||||
if not cb(fd.AsyncFD):
|
if not cb(fd.AsyncFD):
|
||||||
newList.add(cb)
|
newList.add(cb)
|
||||||
|
|
||||||
withData(p.selector, ident.int, adata) do:
|
withData(p.selector, fd.int, adata) do:
|
||||||
# descriptor still present in queue.
|
# descriptor still present in queue.
|
||||||
adata.readList = newList & adata.readList
|
adata.readList = newList & adata.readList
|
||||||
if len(adata.readList) == 0:
|
if len(adata.readList) == 0:
|
||||||
|
|
@ -1308,10 +1309,6 @@ else:
|
||||||
|
|
||||||
proc runOnce(timeout = 500): bool =
|
proc runOnce(timeout = 500): bool =
|
||||||
let p = getGlobalDispatcher()
|
let p = getGlobalDispatcher()
|
||||||
when ioselSupportedPlatform:
|
|
||||||
let customSet = {Event.Timer, Event.Signal, Event.Process,
|
|
||||||
Event.Vnode}
|
|
||||||
|
|
||||||
if p.selector.isEmpty() and p.timers.len == 0 and p.callbacks.len == 0:
|
if p.selector.isEmpty() 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.")
|
||||||
|
|
@ -1346,9 +1343,11 @@ else:
|
||||||
result = true
|
result = true
|
||||||
|
|
||||||
when ioselSupportedPlatform:
|
when ioselSupportedPlatform:
|
||||||
|
const customSet = {Event.Timer, Event.Signal, Event.Process,
|
||||||
|
Event.Vnode}
|
||||||
if (customSet * events) != {}:
|
if (customSet * events) != {}:
|
||||||
isCustomEvent = true
|
isCustomEvent = true
|
||||||
processCustomCallbacks(fd)
|
processCustomCallbacks(p, fd)
|
||||||
result = true
|
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
|
||||||
|
|
@ -1612,7 +1611,7 @@ proc drain*(timeout = 500) =
|
||||||
var curTimeout = timeout
|
var curTimeout = timeout
|
||||||
let start = now()
|
let start = now()
|
||||||
while hasPendingOperations():
|
while hasPendingOperations():
|
||||||
discard runOnce(curTimeout)
|
discard runOnce(curTimeout)
|
||||||
curTimeout -= (now() - start).inMilliseconds.int
|
curTimeout -= (now() - start).inMilliseconds.int
|
||||||
if curTimeout < 0:
|
if curTimeout < 0:
|
||||||
break
|
break
|
||||||
|
|
|
||||||
|
|
@ -157,33 +157,15 @@ proc checkFinished[T](future: Future[T]) =
|
||||||
raise err
|
raise err
|
||||||
|
|
||||||
proc call(callbacks: var CallbackList) =
|
proc call(callbacks: var CallbackList) =
|
||||||
when not defined(nimV2):
|
var current = callbacks
|
||||||
# strictly speaking a little code duplication here, but we strive
|
while true:
|
||||||
# to minimize regressions and I'm not sure I got the 'nimV2' logic
|
if not current.function.isNil:
|
||||||
# right:
|
callSoon(current.function)
|
||||||
var current = callbacks
|
|
||||||
while true:
|
|
||||||
if not current.function.isNil:
|
|
||||||
callSoon(current.function)
|
|
||||||
|
|
||||||
if current.next.isNil:
|
|
||||||
break
|
|
||||||
else:
|
|
||||||
current = current.next[]
|
|
||||||
else:
|
|
||||||
var currentFunc = unown callbacks.function
|
|
||||||
var currentNext = unown callbacks.next
|
|
||||||
|
|
||||||
while true:
|
|
||||||
if not currentFunc.isNil:
|
|
||||||
callSoon(currentFunc)
|
|
||||||
|
|
||||||
if currentNext.isNil:
|
|
||||||
break
|
|
||||||
else:
|
|
||||||
currentFunc = currentNext.function
|
|
||||||
currentNext = unown currentNext.next
|
|
||||||
|
|
||||||
|
if current.next.isNil:
|
||||||
|
break
|
||||||
|
else:
|
||||||
|
current = current.next[]
|
||||||
# callback will be called only once, let GC collect them now
|
# callback will be called only once, let GC collect them now
|
||||||
callbacks.next = nil
|
callbacks.next = nil
|
||||||
callbacks.function = nil
|
callbacks.function = nil
|
||||||
|
|
|
||||||
|
|
@ -54,7 +54,8 @@ proc `callback=`*[T](future: FutureStream[T],
|
||||||
##
|
##
|
||||||
## If the future stream already has data or is finished then ``cb`` will be
|
## If the future stream already has data or is finished then ``cb`` will be
|
||||||
## called immediately.
|
## called immediately.
|
||||||
future.cb = proc () = cb(future)
|
proc named() = cb(future)
|
||||||
|
future.cb = named
|
||||||
if future.queue.len > 0 or future.finished:
|
if future.queue.len > 0 or future.finished:
|
||||||
callSoon(future.cb)
|
callSoon(future.cb)
|
||||||
|
|
||||||
|
|
@ -90,27 +91,26 @@ proc read*[T](future: FutureStream[T]): owned(Future[(bool, T)]) =
|
||||||
## ``FutureStream``.
|
## ``FutureStream``.
|
||||||
var resFut = newFuture[(bool, T)]("FutureStream.take")
|
var resFut = newFuture[(bool, T)]("FutureStream.take")
|
||||||
let savedCb = future.cb
|
let savedCb = future.cb
|
||||||
var newCb =
|
proc newCb(fs: FutureStream[T]) =
|
||||||
proc (fs: FutureStream[T]) =
|
# Exit early if `resFut` is already complete. (See #8994).
|
||||||
# Exit early if `resFut` is already complete. (See #8994).
|
if resFut.finished: return
|
||||||
if resFut.finished: return
|
|
||||||
|
|
||||||
# We don't want this callback called again.
|
# We don't want this callback called again.
|
||||||
#future.cb = nil
|
#future.cb = nil
|
||||||
|
|
||||||
# The return value depends on whether the FutureStream has finished.
|
# The return value depends on whether the FutureStream has finished.
|
||||||
var res: (bool, T)
|
var res: (bool, T)
|
||||||
if finished(fs):
|
if finished(fs):
|
||||||
# Remember, this callback is called when the FutureStream is completed.
|
# Remember, this callback is called when the FutureStream is completed.
|
||||||
res[0] = false
|
res[0] = false
|
||||||
else:
|
else:
|
||||||
res[0] = true
|
res[0] = true
|
||||||
res[1] = fs.queue.popFirst()
|
res[1] = fs.queue.popFirst()
|
||||||
|
|
||||||
resFut.complete(res)
|
resFut.complete(res)
|
||||||
|
|
||||||
# If the saved callback isn't nil then let's call it.
|
# If the saved callback isn't nil then let's call it.
|
||||||
if not savedCb.isNil: savedCb()
|
if not savedCb.isNil: savedCb()
|
||||||
|
|
||||||
if future.queue.len > 0 or future.finished:
|
if future.queue.len > 0 or future.finished:
|
||||||
newCb(future)
|
newCb(future)
|
||||||
|
|
|
||||||
|
|
@ -514,7 +514,7 @@ template withData*[T](s: Selector[T], fd: SocketHandle|int, value,
|
||||||
let fdi = int(fd)
|
let fdi = int(fd)
|
||||||
s.checkFd(fdi)
|
s.checkFd(fdi)
|
||||||
if fdi in s:
|
if fdi in s:
|
||||||
var value = addr(s.getData(fdi))
|
var value = addr(s.fds[fdi].data)
|
||||||
body
|
body
|
||||||
|
|
||||||
template withData*[T](s: Selector[T], fd: SocketHandle|int, value, body1,
|
template withData*[T](s: Selector[T], fd: SocketHandle|int, value, body1,
|
||||||
|
|
@ -523,7 +523,7 @@ template withData*[T](s: Selector[T], fd: SocketHandle|int, value, body1,
|
||||||
let fdi = int(fd)
|
let fdi = int(fd)
|
||||||
s.checkFd(fdi)
|
s.checkFd(fdi)
|
||||||
if fdi in s:
|
if fdi in s:
|
||||||
var value = addr(s.getData(fdi))
|
var value = addr(s.fds[fdi].data)
|
||||||
body1
|
body1
|
||||||
else:
|
else:
|
||||||
body2
|
body2
|
||||||
|
|
|
||||||
|
|
@ -255,7 +255,7 @@ else:
|
||||||
IOSelectorsException* = object of CatchableError
|
IOSelectorsException* = object of CatchableError
|
||||||
|
|
||||||
ReadyKey* = object
|
ReadyKey* = object
|
||||||
fd* : int
|
fd*: int
|
||||||
events*: set[Event]
|
events*: set[Event]
|
||||||
errorCode*: OSErrorCode
|
errorCode*: OSErrorCode
|
||||||
|
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue