FutureStream's cb call behaviour fixed + other fixes.
This commit is contained in:
parent
4a7ea8f865
commit
77071eb767
1 changed files with 6 additions and 1 deletions
|
|
@ -182,7 +182,7 @@ proc `callback=`*[T](future: FutureStream[T],
|
||||||
## If the future stream already has data then ``cb`` will be called
|
## If the future stream already has data then ``cb`` will be called
|
||||||
## immediately.
|
## immediately.
|
||||||
future.cb = proc () = cb(future)
|
future.cb = proc () = cb(future)
|
||||||
if future.queue.len > 0:
|
if future.queue.len > 0 or future.finished:
|
||||||
callSoon(future.cb)
|
callSoon(future.cb)
|
||||||
|
|
||||||
proc injectStacktrace[T](future: Future[T]) =
|
proc injectStacktrace[T](future: Future[T]) =
|
||||||
|
|
@ -257,6 +257,7 @@ proc put*[T](future: FutureStream[T], value: T): Future[void] =
|
||||||
if future.finished:
|
if future.finished:
|
||||||
let msg = "FutureStream is finished and so no longer accepts new data."
|
let msg = "FutureStream is finished and so no longer accepts new data."
|
||||||
result.fail(newException(ValueError, msg))
|
result.fail(newException(ValueError, msg))
|
||||||
|
return
|
||||||
# TODO: Buffering.
|
# TODO: Buffering.
|
||||||
future.queue.enqueue(value)
|
future.queue.enqueue(value)
|
||||||
if not future.cb.isNil: future.cb()
|
if not future.cb.isNil: future.cb()
|
||||||
|
|
@ -294,6 +295,10 @@ proc take*[T](future: FutureStream[T]): Future[(bool, T)] =
|
||||||
if not savedCb.isNil: savedCb()
|
if not savedCb.isNil: savedCb()
|
||||||
return resFut
|
return resFut
|
||||||
|
|
||||||
|
proc len*[T](future: FutureStream[T]): int =
|
||||||
|
## Returns the amount of data pieces inside the stream.
|
||||||
|
future.queue.len
|
||||||
|
|
||||||
proc asyncCheck*[T](future: Future[T]) =
|
proc asyncCheck*[T](future: Future[T]) =
|
||||||
## Sets a callback on ``future`` which raises an exception if the future
|
## Sets a callback on ``future`` which raises an exception if the future
|
||||||
## finished with an error.
|
## finished with an error.
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue