Merge pull request #5041 from cheatfate/upcoming_async_update1
upcoming_asyncdispatch, make it compatible again
This commit is contained in:
commit
7c9a61f241
1 changed files with 60 additions and 36 deletions
|
|
@ -9,7 +9,7 @@
|
||||||
|
|
||||||
include "system/inclrtl"
|
include "system/inclrtl"
|
||||||
|
|
||||||
import os, oids, tables, strutils, times, heapqueue
|
import os, oids, tables, strutils, times, heapqueue, lists
|
||||||
|
|
||||||
import nativesockets, net, queues
|
import nativesockets, net, queues
|
||||||
|
|
||||||
|
|
@ -1095,9 +1095,11 @@ else:
|
||||||
AsyncFD* = distinct cint
|
AsyncFD* = distinct cint
|
||||||
Callback = proc (fd: AsyncFD): bool {.closure,gcsafe.}
|
Callback = proc (fd: AsyncFD): bool {.closure,gcsafe.}
|
||||||
|
|
||||||
|
DoublyLinkedListRef = ref DoublyLinkedList[Callback]
|
||||||
|
|
||||||
AsyncData = object
|
AsyncData = object
|
||||||
readCB: Callback
|
readCBs: DoublyLinkedListRef
|
||||||
writeCB: Callback
|
writeCBs: DoublyLinkedListRef
|
||||||
|
|
||||||
AsyncEvent* = distinct SelectEvent
|
AsyncEvent* = distinct SelectEvent
|
||||||
|
|
||||||
|
|
@ -1121,7 +1123,10 @@ else:
|
||||||
|
|
||||||
proc register*(fd: AsyncFD) =
|
proc register*(fd: AsyncFD) =
|
||||||
let p = getGlobalDispatcher()
|
let p = getGlobalDispatcher()
|
||||||
var data = AsyncData()
|
var data = AsyncData(
|
||||||
|
readCBs: DoublyLinkedListRef(),
|
||||||
|
writeCBs: DoublyLinkedListRef()
|
||||||
|
)
|
||||||
p.selector.registerHandle(fd.SocketHandle, {}, data)
|
p.selector.registerHandle(fd.SocketHandle, {}, data)
|
||||||
|
|
||||||
proc newAsyncNativeSocket*(domain: cint, sockType: cint,
|
proc newAsyncNativeSocket*(domain: cint, sockType: cint,
|
||||||
|
|
@ -1156,8 +1161,9 @@ else:
|
||||||
let p = getGlobalDispatcher()
|
let p = getGlobalDispatcher()
|
||||||
var newEvents = {Event.Read}
|
var newEvents = {Event.Read}
|
||||||
withData(p.selector, fd.SocketHandle, adata) do:
|
withData(p.selector, fd.SocketHandle, adata) do:
|
||||||
adata.readCB = cb
|
adata.readCBs[].append(cb)
|
||||||
if adata.writeCB != nil:
|
newEvents.incl(Event.Read)
|
||||||
|
if not isNil(adata.writeCBs.head):
|
||||||
newEvents.incl(Event.Write)
|
newEvents.incl(Event.Write)
|
||||||
do:
|
do:
|
||||||
raise newException(ValueError, "File descriptor not registered.")
|
raise newException(ValueError, "File descriptor not registered.")
|
||||||
|
|
@ -1167,8 +1173,9 @@ else:
|
||||||
let p = getGlobalDispatcher()
|
let p = getGlobalDispatcher()
|
||||||
var newEvents = {Event.Write}
|
var newEvents = {Event.Write}
|
||||||
withData(p.selector, fd.SocketHandle, adata) do:
|
withData(p.selector, fd.SocketHandle, adata) do:
|
||||||
adata.writeCB = cb
|
adata.writeCBs[].append(cb)
|
||||||
if adata.readCB != nil:
|
newEvents.incl(Event.Write)
|
||||||
|
if not isNil(adata.readCBs.head):
|
||||||
newEvents.incl(Event.Read)
|
newEvents.incl(Event.Read)
|
||||||
do:
|
do:
|
||||||
raise newException(ValueError, "File descriptor not registered.")
|
raise newException(ValueError, "File descriptor not registered.")
|
||||||
|
|
@ -1195,30 +1202,31 @@ else:
|
||||||
let events = keys[i].events
|
let events = keys[i].events
|
||||||
|
|
||||||
if Event.Read in events or events == {Event.Error}:
|
if Event.Read in events or events == {Event.Error}:
|
||||||
let cb = keys[i].data.readCB
|
for node in keys[i].data.readCBs[].nodes():
|
||||||
|
let cb = node.value
|
||||||
if cb != nil:
|
if cb != nil:
|
||||||
if cb(fd.AsyncFD):
|
if cb(fd.AsyncFD):
|
||||||
p.selector.withData(fd, adata) do:
|
keys[i].data.readCBs[].remove(node)
|
||||||
if adata.readCB == cb:
|
else:
|
||||||
adata.readCB = nil
|
break
|
||||||
|
|
||||||
if Event.Write in events or events == {Event.Error}:
|
if Event.Write in events or events == {Event.Error}:
|
||||||
let cb = keys[i].data.writeCB
|
for node in keys[i].data.writeCBs[].nodes():
|
||||||
|
let cb = node.value
|
||||||
if cb != nil:
|
if cb != nil:
|
||||||
if cb(fd.AsyncFD):
|
if cb(fd.AsyncFD):
|
||||||
p.selector.withData(fd, adata) do:
|
keys[i].data.writeCBs[].remove(node)
|
||||||
if adata.writeCB == cb:
|
else:
|
||||||
adata.writeCB = nil
|
break
|
||||||
|
|
||||||
when supportedPlatform:
|
when supportedPlatform:
|
||||||
if (customSet * events) != {}:
|
if (customSet * events) != {}:
|
||||||
let cb = keys[i].data.readCB
|
for node in keys[i].data.readCBs[].nodes():
|
||||||
|
let cb = node.value
|
||||||
doAssert(cb != nil)
|
doAssert(cb != nil)
|
||||||
custom = true
|
custom = true
|
||||||
if cb(fd.AsyncFD):
|
if cb(fd.AsyncFD):
|
||||||
p.selector.withData(fd, adata) do:
|
keys[i].data.readCBs[].remove(node)
|
||||||
if adata.readCB == cb:
|
|
||||||
adata.readCB = nil
|
|
||||||
p.selector.unregister(fd)
|
p.selector.unregister(fd)
|
||||||
|
|
||||||
# because state `data` can be modified in callback we need to update
|
# because state `data` can be modified in callback we need to update
|
||||||
|
|
@ -1227,8 +1235,8 @@ else:
|
||||||
var update = false
|
var update = false
|
||||||
var newEvents: set[Event] = {}
|
var newEvents: set[Event] = {}
|
||||||
p.selector.withData(fd, adata) do:
|
p.selector.withData(fd, adata) do:
|
||||||
if adata.readCB != nil: incl(newEvents, Event.Read)
|
if not isNil(adata.readCBs.head): incl(newEvents, Event.Read)
|
||||||
if adata.writeCB != nil: incl(newEvents, Event.Write)
|
if not isNil(adata.writeCBs.head): incl(newEvents, Event.Write)
|
||||||
update = true
|
update = true
|
||||||
if update:
|
if update:
|
||||||
p.selector.updateHandle(fd, newEvents)
|
p.selector.updateHandle(fd, newEvents)
|
||||||
|
|
@ -1491,21 +1499,33 @@ else:
|
||||||
## ``oneshot`` - if ``true`` only one event will be dispatched,
|
## ``oneshot`` - if ``true`` only one event will be dispatched,
|
||||||
## if ``false`` continuous events every ``timeout`` milliseconds.
|
## if ``false`` continuous events every ``timeout`` milliseconds.
|
||||||
let p = getGlobalDispatcher()
|
let p = getGlobalDispatcher()
|
||||||
var data = AsyncData(readCB: cb)
|
var data = AsyncData(
|
||||||
|
readCBs: DoublyLinkedListRef(),
|
||||||
|
writeCBs: DoublyLinkedListRef()
|
||||||
|
)
|
||||||
|
data.readCBs[].append(cb)
|
||||||
p.selector.registerTimer(timeout, oneshot, data)
|
p.selector.registerTimer(timeout, oneshot, data)
|
||||||
|
|
||||||
proc addSignal*(signal: int, cb: Callback) =
|
proc addSignal*(signal: int, cb: Callback) =
|
||||||
## Start watching signal ``signal``, and when signal appears, call the
|
## Start watching signal ``signal``, and when signal appears, call the
|
||||||
## callback ``cb``.
|
## callback ``cb``.
|
||||||
let p = getGlobalDispatcher()
|
let p = getGlobalDispatcher()
|
||||||
var data = AsyncData(readCB: cb)
|
var data = AsyncData(
|
||||||
|
readCBs: DoublyLinkedListRef(),
|
||||||
|
writeCBs: DoublyLinkedListRef()
|
||||||
|
)
|
||||||
|
data.readCBs[].append(cb)
|
||||||
p.selector.registerSignal(signal, data)
|
p.selector.registerSignal(signal, data)
|
||||||
|
|
||||||
proc addProcess*(pid: int, cb: Callback) =
|
proc addProcess*(pid: int, cb: Callback) =
|
||||||
## Start watching for process exit with pid ``pid``, and then call
|
## Start watching for process exit with pid ``pid``, and then call
|
||||||
## the callback ``cb``.
|
## the callback ``cb``.
|
||||||
let p = getGlobalDispatcher()
|
let p = getGlobalDispatcher()
|
||||||
var data = AsyncData(readCB: cb)
|
var data = AsyncData(
|
||||||
|
readCBs: DoublyLinkedListRef(),
|
||||||
|
writeCBs: DoublyLinkedListRef()
|
||||||
|
)
|
||||||
|
data.readCBs[].append(cb)
|
||||||
p.selector.registerProcess(pid, data)
|
p.selector.registerProcess(pid, data)
|
||||||
|
|
||||||
proc newAsyncEvent*(): AsyncEvent =
|
proc newAsyncEvent*(): AsyncEvent =
|
||||||
|
|
@ -1524,7 +1544,11 @@ else:
|
||||||
## Start watching for event ``ev``, and call callback ``cb``, when
|
## Start watching for event ``ev``, and call callback ``cb``, when
|
||||||
## ev will be set to signaled state.
|
## ev will be set to signaled state.
|
||||||
let p = getGlobalDispatcher()
|
let p = getGlobalDispatcher()
|
||||||
var data = AsyncData(readCB: cb)
|
var data = AsyncData(
|
||||||
|
readCBs: DoublyLinkedListRef(),
|
||||||
|
writeCBs: DoublyLinkedListRef()
|
||||||
|
)
|
||||||
|
data.readCBs[].append(cb)
|
||||||
p.selector.registerEvent(SelectEvent(ev), data)
|
p.selector.registerEvent(SelectEvent(ev), data)
|
||||||
|
|
||||||
proc sleepAsync*(ms: int): Future[void] =
|
proc sleepAsync*(ms: int): Future[void] =
|
||||||
|
|
@ -1591,7 +1615,7 @@ proc recvLine*(socket: AsyncFD): Future[string] {.async.} =
|
||||||
## **Note**: This procedure is mostly used for testing. You likely want to
|
## **Note**: This procedure is mostly used for testing. You likely want to
|
||||||
## use ``asyncnet.recvLine`` instead.
|
## use ``asyncnet.recvLine`` instead.
|
||||||
|
|
||||||
template addNLIfEmpty(): stmt =
|
template addNLIfEmpty(): typed =
|
||||||
if result.len == 0:
|
if result.len == 0:
|
||||||
result.add("\c\L")
|
result.add("\c\L")
|
||||||
|
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue