asyncdispatch: split asyncfutures into its own module
This slightly changes behaviour of callSoon - before loop is initialized, callSoon will call the function immediately.
This commit is contained in:
parent
9e12db4459
commit
e86863e8f5
4 changed files with 52 additions and 22 deletions
|
|
@ -9,11 +9,11 @@
|
||||||
|
|
||||||
include "system/inclrtl"
|
include "system/inclrtl"
|
||||||
|
|
||||||
import os, tables, strutils, times, heapqueue, options
|
import os, tables, strutils, times, heapqueue, options, asyncfutures
|
||||||
|
|
||||||
import nativesockets, net, deques
|
import nativesockets, net, deques
|
||||||
|
|
||||||
export Port, SocketFlag
|
export Port, SocketFlag
|
||||||
|
export asyncfutures
|
||||||
|
|
||||||
#{.injectStmt: newGcInvariant().}
|
#{.injectStmt: newGcInvariant().}
|
||||||
|
|
||||||
|
|
@ -159,8 +159,6 @@ export Port, SocketFlag
|
||||||
|
|
||||||
# TODO: Check if yielded future is nil and throw a more meaningful exception
|
# TODO: Check if yielded future is nil and throw a more meaningful exception
|
||||||
|
|
||||||
include includes/asyncfutures
|
|
||||||
|
|
||||||
type
|
type
|
||||||
PDispatcherBase = ref object of RootRef
|
PDispatcherBase = ref object of RootRef
|
||||||
timers*: HeapQueue[tuple[finishAt: float, fut: Future[void]]]
|
timers*: HeapQueue[tuple[finishAt: float, fut: Future[void]]]
|
||||||
|
|
@ -190,6 +188,12 @@ proc adjustedTimeout(p: PDispatcherBase, timeout: int): int {.inline.} =
|
||||||
result = int((timerTimeout - curTime) * 1000)
|
result = int((timerTimeout - curTime) * 1000)
|
||||||
if result < 0: result = 0
|
if result < 0: result = 0
|
||||||
|
|
||||||
|
proc callSoon*(cbproc: proc ()) {.gcsafe.}
|
||||||
|
|
||||||
|
proc initGlobalDispatcher =
|
||||||
|
if asyncfutures.callSoonProc == nil:
|
||||||
|
asyncfutures.callSoonProc = callSoon
|
||||||
|
|
||||||
when defined(windows) or defined(nimdoc):
|
when defined(windows) or defined(nimdoc):
|
||||||
import winlean, sets, hashes
|
import winlean, sets, hashes
|
||||||
type
|
type
|
||||||
|
|
@ -237,15 +241,17 @@ when defined(windows) or defined(nimdoc):
|
||||||
result.callbacks = initDeque[proc ()](64)
|
result.callbacks = initDeque[proc ()](64)
|
||||||
|
|
||||||
var gDisp{.threadvar.}: PDispatcher ## Global dispatcher
|
var gDisp{.threadvar.}: PDispatcher ## Global dispatcher
|
||||||
proc getGlobalDispatcher*(): PDispatcher =
|
|
||||||
## Retrieves the global thread-local dispatcher.
|
|
||||||
if gDisp.isNil: gDisp = newDispatcher()
|
|
||||||
result = gDisp
|
|
||||||
|
|
||||||
proc setGlobalDispatcher*(disp: PDispatcher) =
|
proc setGlobalDispatcher*(disp: PDispatcher) =
|
||||||
if not gDisp.isNil:
|
if not gDisp.isNil:
|
||||||
assert gDisp.callbacks.len == 0
|
assert gDisp.callbacks.len == 0
|
||||||
gDisp = disp
|
gDisp = disp
|
||||||
|
initGlobalDispatcher()
|
||||||
|
|
||||||
|
proc getGlobalDispatcher*(): PDispatcher =
|
||||||
|
if gDisp.isNil:
|
||||||
|
setGlobalDispatcher(newDispatcher())
|
||||||
|
result = gDisp
|
||||||
|
|
||||||
proc register*(fd: AsyncFD) =
|
proc register*(fd: AsyncFD) =
|
||||||
## Registers ``fd`` with the dispatcher.
|
## Registers ``fd`` with the dispatcher.
|
||||||
|
|
@ -932,14 +938,17 @@ else:
|
||||||
result.callbacks = initDeque[proc ()](64)
|
result.callbacks = initDeque[proc ()](64)
|
||||||
|
|
||||||
var gDisp{.threadvar.}: PDispatcher ## Global dispatcher
|
var gDisp{.threadvar.}: PDispatcher ## Global dispatcher
|
||||||
proc getGlobalDispatcher*(): PDispatcher =
|
|
||||||
if gDisp.isNil: gDisp = newDispatcher()
|
|
||||||
result = gDisp
|
|
||||||
|
|
||||||
proc setGlobalDispatcher*(disp: PDispatcher) =
|
proc setGlobalDispatcher*(disp: PDispatcher) =
|
||||||
if not gDisp.isNil:
|
if not gDisp.isNil:
|
||||||
assert gDisp.callbacks.len == 0
|
assert gDisp.callbacks.len == 0
|
||||||
gDisp = disp
|
gDisp = disp
|
||||||
|
initGlobalDispatcher()
|
||||||
|
|
||||||
|
proc getGlobalDispatcher*(): PDispatcher =
|
||||||
|
if gDisp.isNil:
|
||||||
|
setGlobalDispatcher(newDispatcher())
|
||||||
|
result = gDisp
|
||||||
|
|
||||||
proc update(fd: AsyncFD, events: set[Event]) =
|
proc update(fd: AsyncFD, events: set[Event]) =
|
||||||
let p = getGlobalDispatcher()
|
let p = getGlobalDispatcher()
|
||||||
|
|
|
||||||
|
|
@ -1,3 +1,6 @@
|
||||||
|
include "system/inclrtl"
|
||||||
|
|
||||||
|
import os, tables, strutils, times, heapqueue, options, deques
|
||||||
|
|
||||||
# TODO: This shouldn't need to be included, but should ideally be exported.
|
# TODO: This shouldn't need to be included, but should ideally be exported.
|
||||||
type
|
type
|
||||||
|
|
@ -30,7 +33,14 @@ type
|
||||||
when not defined(release):
|
when not defined(release):
|
||||||
var currentID = 0
|
var currentID = 0
|
||||||
|
|
||||||
proc callSoon*(cbproc: proc ()) {.gcsafe.}
|
var callSoonProc* {.threadvar.}: (proc(cbproc: proc ()) {.gcsafe.})
|
||||||
|
|
||||||
|
proc callSoon(cbproc: proc ()) =
|
||||||
|
if callSoonProc == nil:
|
||||||
|
# Loop not initialized yet. Call the function directly to allow setup code to use futures.
|
||||||
|
cbproc()
|
||||||
|
else:
|
||||||
|
callSoonProc(cbproc)
|
||||||
|
|
||||||
template setupFutureBase(fromProc: string) =
|
template setupFutureBase(fromProc: string) =
|
||||||
new(result)
|
new(result)
|
||||||
|
|
@ -9,11 +9,12 @@
|
||||||
|
|
||||||
include "system/inclrtl"
|
include "system/inclrtl"
|
||||||
|
|
||||||
import os, tables, strutils, times, heapqueue, lists, options
|
import os, tables, strutils, times, heapqueue, lists, options, asyncfutures
|
||||||
|
|
||||||
import nativesockets, net, deques
|
import nativesockets, net, deques
|
||||||
|
|
||||||
export Port, SocketFlag
|
export Port, SocketFlag
|
||||||
|
export asyncfutures
|
||||||
|
|
||||||
#{.injectStmt: newGcInvariant().}
|
#{.injectStmt: newGcInvariant().}
|
||||||
|
|
||||||
|
|
@ -130,8 +131,6 @@ export Port, SocketFlag
|
||||||
|
|
||||||
# TODO: Check if yielded future is nil and throw a more meaningful exception
|
# TODO: Check if yielded future is nil and throw a more meaningful exception
|
||||||
|
|
||||||
include "../includes/asyncfutures"
|
|
||||||
|
|
||||||
type
|
type
|
||||||
PDispatcherBase = ref object of RootRef
|
PDispatcherBase = ref object of RootRef
|
||||||
timers: HeapQueue[tuple[finishAt: float, fut: Future[void]]]
|
timers: HeapQueue[tuple[finishAt: float, fut: Future[void]]]
|
||||||
|
|
@ -161,6 +160,12 @@ proc adjustedTimeout(p: PDispatcherBase, timeout: int): int {.inline.} =
|
||||||
result = int((timerTimeout - curTime) * 1000)
|
result = int((timerTimeout - curTime) * 1000)
|
||||||
if result < 0: result = 0
|
if result < 0: result = 0
|
||||||
|
|
||||||
|
proc callSoon*(cbproc: proc ()) {.gcsafe.}
|
||||||
|
|
||||||
|
proc initGlobalDispatcher =
|
||||||
|
if asyncfutures.callSoonProc == nil:
|
||||||
|
asyncfutures.callSoonProc = callSoon
|
||||||
|
|
||||||
when defined(windows) or defined(nimdoc):
|
when defined(windows) or defined(nimdoc):
|
||||||
import winlean, sets, hashes
|
import winlean, sets, hashes
|
||||||
type
|
type
|
||||||
|
|
@ -214,15 +219,17 @@ when defined(windows) or defined(nimdoc):
|
||||||
result.callbacks = initDeque[proc ()](64)
|
result.callbacks = initDeque[proc ()](64)
|
||||||
|
|
||||||
var gDisp{.threadvar.}: PDispatcher ## Global dispatcher
|
var gDisp{.threadvar.}: PDispatcher ## Global dispatcher
|
||||||
proc getGlobalDispatcher*(): PDispatcher =
|
|
||||||
## Retrieves the global thread-local dispatcher.
|
|
||||||
if gDisp.isNil: gDisp = newDispatcher()
|
|
||||||
result = gDisp
|
|
||||||
|
|
||||||
proc setGlobalDispatcher*(disp: PDispatcher) =
|
proc setGlobalDispatcher*(disp: PDispatcher) =
|
||||||
if not gDisp.isNil:
|
if not gDisp.isNil:
|
||||||
assert gDisp.callbacks.len == 0
|
assert gDisp.callbacks.len == 0
|
||||||
gDisp = disp
|
gDisp = disp
|
||||||
|
initGlobalDispatcher()
|
||||||
|
|
||||||
|
proc getGlobalDispatcher*(): PDispatcher =
|
||||||
|
if gDisp.isNil:
|
||||||
|
setGlobalDispatcher(newDispatcher())
|
||||||
|
result = gDisp
|
||||||
|
|
||||||
proc register*(fd: AsyncFD) =
|
proc register*(fd: AsyncFD) =
|
||||||
## Registers ``fd`` with the dispatcher.
|
## Registers ``fd`` with the dispatcher.
|
||||||
|
|
@ -1081,14 +1088,17 @@ else:
|
||||||
result.callbacks = initDeque[proc ()](InitDelayedCallbackListSize)
|
result.callbacks = initDeque[proc ()](InitDelayedCallbackListSize)
|
||||||
|
|
||||||
var gDisp{.threadvar.}: PDispatcher ## Global dispatcher
|
var gDisp{.threadvar.}: PDispatcher ## Global dispatcher
|
||||||
proc getGlobalDispatcher*(): PDispatcher =
|
|
||||||
if gDisp.isNil: gDisp = newDispatcher()
|
|
||||||
result = gDisp
|
|
||||||
|
|
||||||
proc setGlobalDispatcher*(disp: PDispatcher) =
|
proc setGlobalDispatcher*(disp: PDispatcher) =
|
||||||
if not gDisp.isNil:
|
if not gDisp.isNil:
|
||||||
assert gDisp.callbacks.len == 0
|
assert gDisp.callbacks.len == 0
|
||||||
gDisp = disp
|
gDisp = disp
|
||||||
|
initGlobalDispatcher()
|
||||||
|
|
||||||
|
proc getGlobalDispatcher*(): PDispatcher =
|
||||||
|
if gDisp.isNil:
|
||||||
|
setGlobalDispatcher(newDispatcher())
|
||||||
|
result = gDisp
|
||||||
|
|
||||||
proc register*(fd: AsyncFD) =
|
proc register*(fd: AsyncFD) =
|
||||||
let p = getGlobalDispatcher()
|
let p = getGlobalDispatcher()
|
||||||
|
|
|
||||||
|
|
@ -17,5 +17,6 @@ proc asyncRecursionTest*(): Future[int] {.async.} =
|
||||||
inc(i)
|
inc(i)
|
||||||
|
|
||||||
when isMainModule:
|
when isMainModule:
|
||||||
|
setGlobalDispatcher(newDispatcher())
|
||||||
var i = waitFor asyncRecursionTest()
|
var i = waitFor asyncRecursionTest()
|
||||||
echo i
|
echo i
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue