move future to method

This commit is contained in:
Jaremy Creechley 2023-09-14 15:38:02 -07:00
parent 343897dab5
commit 911abb43aa
No known key found for this signature in database
GPG Key ID: 4E66FB67B21D3300
4 changed files with 49 additions and 60 deletions

View File

@ -81,16 +81,51 @@ method get*(
return ret.convert(seq[byte])
import ./threads/then
import std/os
method put*(
self: ThreadProxyDatastore,
key: Key,
data: seq[byte]
): Future[?!void] {.async.} =
): Future[?!void] =
echoed "put new request thr: ", $getThreadId()
echoed "put request args: ", $getThreadId()
let res = await put(self.tds, key, data)
return res
let tds = self.tds
var putRes = newFuture[?!void]("threadbackend.put(tds, key, data)")
let sig = SharedSignal.new(2)
acquireSig(sig).
then(proc () =
echoed "got tresFut"
let
ret = newSharedPtr(ThreadResult[void])
bkey = KeyBuffer.new(key)
bval = DataBuffer.new(data)
tds[].tp.spawn putTask(sig, ret, tds, bkey, bval)
wait(sig).
then(proc () =
sig.decr()
os.sleep(400)
var ret = ret
let val = ret.convert(void)
putRes.complete(val)
).cancelled(proc() =
sig.decr()
).catch(proc(e: ref CatchableError) =
sig.decr()
doAssert false, "will not be triggered"
)
).catch(proc(e: ref CatchableError) =
var res: ?!void
res.err(e)
putRes.complete(res)
)
return putRes
method put*(
self: ThreadProxyDatastore,

View File

@ -150,58 +150,6 @@ proc putTask*(
sig.decr()
echoed "putTask: FINISH\n"
import then
proc put*(
tds: ThreadDatastorePtr,
key: Key,
data: seq[byte]
): Future[?!void] =
echoed "put request args: ", $getThreadId()
var putRes = newFuture[?!void]("threadbackend.put(tds, key, data)")
let sigFut = SharedSignal.new(2)
sigFut.
then(proc (sig: SharedSignal) =
echoed "got tresFut"
let
ret = newSharedPtr(ThreadResult[void])
bkey = KeyBuffer.new(key)
bval = DataBuffer.new(data)
echoed "spawn put request: ", $getThreadId()
# this spawns the taskpool Task
# but we can't wait on it directly - we use wait(ret[].sig)
echo "\n"
tds[].tp.spawn putTask(sig, ret, tds, bkey, bval)
wait(sig).
then(proc () =
sig.decr()
echo "\n"
os.sleep(400)
echoed "put request done "
var ret = ret
let val = ret.convert(void)
putRes.complete(val)
).cancelled(proc() =
sig.decr()
echoed "put request cancelled "
discard
).catch(proc(e: ref CatchableError) =
sig.decr()
doAssert false, "will not be triggered"
)
).catch(proc(e: ref CatchableError) =
echoed "err tresFut"
var res: ?!void
res.err(e)
putRes.complete(res)
)
return putRes
proc deleteTask*(

View File

@ -58,7 +58,7 @@ proc newThreadResult*[T](
{.error: "only thread safe types can be used".}
let res = newSharedPtr(ThreadResult[T])
res[].sig = await SharedSignal.new(0)
res[].sig = await SharedSignal.new()
res
proc `=destroy`*[T](res: var ThreadResult[T]) {.raises: [].} =

View File

@ -90,11 +90,17 @@ proc `=destroy`*[T](x: var SharedSignalObj) =
release(x.sigptr)
x.sigptr = nil
proc new*(tp: typedesc[SharedSignal],
count: int): Future[SharedSignal] {.async.} =
result = newSharedPtr[SharedSignalObj](SharedSignalObj, manualCount = count)
proc new*(tp: typedesc[SharedSignal]): Future[SharedSignal] {.async.} =
result = newSharedPtr[SharedSignalObj](SharedSignalObj)
result[].sigptr = await getThreadSignal()
proc new*(tp: typedesc[SharedSignal],
count: int): SharedSignal =
result = newSharedPtr[SharedSignalObj](SharedSignalObj, manualCount = count)
proc acquireSig*(sig: SharedSignal): Future[void] {.async.} =
sig[].sigptr = await getThreadSignal()
proc wait*(sig: SharedSignal): Future[void] =
sig[].sigptr.wait()