Future: support for multiple callbacks
This commit is contained in:
parent
93827e6ab8
commit
797690ba3f
2 changed files with 82 additions and 21 deletions
|
|
@ -4,8 +4,15 @@ 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
|
||||||
|
CallbackFunc = proc () {.closure, gcsafe.}
|
||||||
|
|
||||||
|
CallbackList = object
|
||||||
|
function: CallbackFunc
|
||||||
|
next: ref CallbackList
|
||||||
|
|
||||||
FutureBase* = ref object of RootObj ## Untyped future.
|
FutureBase* = ref object of RootObj ## Untyped future.
|
||||||
cb: proc () {.closure,gcsafe.}
|
callbacks: CallbackList
|
||||||
|
|
||||||
finished: bool
|
finished: bool
|
||||||
error*: ref Exception ## Stored exception
|
error*: ref Exception ## Stored exception
|
||||||
errorStackTrace*: string
|
errorStackTrace*: string
|
||||||
|
|
@ -86,6 +93,33 @@ proc checkFinished[T](future: Future[T]) =
|
||||||
err.cause = future
|
err.cause = future
|
||||||
raise err
|
raise err
|
||||||
|
|
||||||
|
proc call(callbacks: var CallbackList) =
|
||||||
|
var current = callbacks
|
||||||
|
|
||||||
|
while true:
|
||||||
|
if current.function != nil:
|
||||||
|
callSoon(current.function)
|
||||||
|
|
||||||
|
if current.next == nil:
|
||||||
|
break
|
||||||
|
else:
|
||||||
|
current = current.next[]
|
||||||
|
|
||||||
|
# callback will be called only once, let GC collect them now
|
||||||
|
callbacks.next = nil
|
||||||
|
callbacks.function = nil
|
||||||
|
|
||||||
|
proc add(callbacks: var CallbackList, function: CallbackFunc) =
|
||||||
|
if callbacks.function == nil:
|
||||||
|
callbacks.function = function
|
||||||
|
assert callbacks.next == nil
|
||||||
|
else:
|
||||||
|
let newNext = new(ref CallbackList)
|
||||||
|
newNext.function = callbacks.function
|
||||||
|
newNext.next = callbacks.next
|
||||||
|
callbacks.next = newNext
|
||||||
|
callbacks.function = function
|
||||||
|
|
||||||
proc complete*[T](future: Future[T], val: T) =
|
proc complete*[T](future: Future[T], val: T) =
|
||||||
## Completes ``future`` with value ``val``.
|
## Completes ``future`` with value ``val``.
|
||||||
#assert(not future.finished, "Future already finished, cannot finish twice.")
|
#assert(not future.finished, "Future already finished, cannot finish twice.")
|
||||||
|
|
@ -93,8 +127,7 @@ proc complete*[T](future: Future[T], val: T) =
|
||||||
assert(future.error == nil)
|
assert(future.error == nil)
|
||||||
future.value = val
|
future.value = val
|
||||||
future.finished = true
|
future.finished = true
|
||||||
if future.cb != nil:
|
future.callbacks.call()
|
||||||
future.cb()
|
|
||||||
|
|
||||||
proc complete*(future: Future[void]) =
|
proc complete*(future: Future[void]) =
|
||||||
## Completes a void ``future``.
|
## Completes a void ``future``.
|
||||||
|
|
@ -102,8 +135,7 @@ proc complete*(future: Future[void]) =
|
||||||
checkFinished(future)
|
checkFinished(future)
|
||||||
assert(future.error == nil)
|
assert(future.error == nil)
|
||||||
future.finished = true
|
future.finished = true
|
||||||
if future.cb != nil:
|
future.callbacks.call()
|
||||||
future.cb()
|
|
||||||
|
|
||||||
proc complete*[T](future: FutureVar[T]) =
|
proc complete*[T](future: FutureVar[T]) =
|
||||||
## Completes a ``FutureVar``.
|
## Completes a ``FutureVar``.
|
||||||
|
|
@ -111,8 +143,7 @@ proc complete*[T](future: FutureVar[T]) =
|
||||||
checkFinished(fut)
|
checkFinished(fut)
|
||||||
assert(fut.error == nil)
|
assert(fut.error == nil)
|
||||||
fut.finished = true
|
fut.finished = true
|
||||||
if fut.cb != nil:
|
fut.callbacks.call
|
||||||
fut.cb()
|
|
||||||
|
|
||||||
proc complete*[T](future: FutureVar[T], val: T) =
|
proc complete*[T](future: FutureVar[T], val: T) =
|
||||||
## Completes a ``FutureVar`` with value ``val``.
|
## Completes a ``FutureVar`` with value ``val``.
|
||||||
|
|
@ -134,26 +165,36 @@ proc fail*[T](future: Future[T], error: ref Exception) =
|
||||||
future.error = error
|
future.error = error
|
||||||
future.errorStackTrace =
|
future.errorStackTrace =
|
||||||
if getStackTrace(error) == "": getStackTrace() else: getStackTrace(error)
|
if getStackTrace(error) == "": getStackTrace() else: getStackTrace(error)
|
||||||
if future.cb != nil:
|
future.callbacks.call
|
||||||
future.cb()
|
|
||||||
|
proc clearCallbacks(future: FutureBase) =
|
||||||
|
future.callbacks.function = nil
|
||||||
|
future.callbacks.next = nil
|
||||||
|
|
||||||
|
proc addCallback*(future: FutureBase, cb: proc() {.closure,gcsafe.}) =
|
||||||
|
## Adds the callbacks proc to be called when the future completes.
|
||||||
|
##
|
||||||
|
## If future has already completed then ``cb`` will be called immediately.
|
||||||
|
assert cb != nil
|
||||||
|
if future.finished:
|
||||||
|
callSoon(cb)
|
||||||
else:
|
else:
|
||||||
# This is to prevent exceptions from being silently ignored when a future
|
future.callbacks.add cb
|
||||||
# is discarded.
|
|
||||||
# TODO: This may turn out to be a bad idea.
|
proc addCallback*[T](future: Future[T], cb: proc(future: Future[T]) {.closure,gcsafe.}) =
|
||||||
# Turns out this is a bad idea.
|
## Adds the callbacks proc to be called when the future completes.
|
||||||
#raise error
|
##
|
||||||
discard
|
## If future has already completed then ``cb`` will be called immediately.
|
||||||
|
future.addCallback proc() = cb(future)
|
||||||
|
|
||||||
proc `callback=`*(future: FutureBase, cb: proc () {.closure,gcsafe.}) =
|
proc `callback=`*(future: FutureBase, cb: proc () {.closure,gcsafe.}) =
|
||||||
## Sets the callback proc to be called when the future completes.
|
## Clears the list of callbacks and sets the callback proc to be called when the future completes.
|
||||||
##
|
##
|
||||||
## If future has already completed then ``cb`` will be called immediately.
|
## If future has already completed then ``cb`` will be called immediately.
|
||||||
##
|
##
|
||||||
## **Note**: You most likely want the other ``callback`` setter which
|
## It's recommended to use ``addCallback`` or ``then`` instead.
|
||||||
## passes ``future`` as a param to the callback.
|
future.clearCallbacks
|
||||||
future.cb = cb
|
future.addCallback cb
|
||||||
if future.finished:
|
|
||||||
callSoon(future.cb)
|
|
||||||
|
|
||||||
proc `callback=`*[T](future: Future[T],
|
proc `callback=`*[T](future: Future[T],
|
||||||
cb: proc (future: Future[T]) {.closure,gcsafe.}) =
|
cb: proc (future: Future[T]) {.closure,gcsafe.}) =
|
||||||
|
|
|
||||||
20
tests/async/tcallbacks.nim
Normal file
20
tests/async/tcallbacks.nim
Normal file
|
|
@ -0,0 +1,20 @@
|
||||||
|
discard """
|
||||||
|
exitcode: 0
|
||||||
|
output: '''3
|
||||||
|
2
|
||||||
|
1
|
||||||
|
5
|
||||||
|
'''
|
||||||
|
"""
|
||||||
|
import asyncfutures
|
||||||
|
|
||||||
|
let f1: Future[int] = newFuture[int]()
|
||||||
|
f1.addCallback(proc() = echo 1)
|
||||||
|
f1.addCallback(proc() = echo 2)
|
||||||
|
f1.addCallback(proc() = echo 3)
|
||||||
|
f1.complete(10)
|
||||||
|
|
||||||
|
let f2: Future[int] = newFuture[int]()
|
||||||
|
f2.addCallback(proc() = echo 4)
|
||||||
|
f2.callback = proc() = echo 5
|
||||||
|
f2.complete(10)
|
||||||
Loading…
Add table
Add a link
Reference in a new issue