From 911abb43aaf0c89101f0a0b27d95a6b17622f956 Mon Sep 17 00:00:00 2001 From: Jaremy Creechley Date: Thu, 14 Sep 2023 15:38:02 -0700 Subject: [PATCH] move future to method --- datastore/threadproxyds.nim | 43 +++++++++++++++++++-- datastore/threads/threadbackend.nim | 52 -------------------------- datastore/threads/threadresults.nim | 2 +- datastore/threads/threadsignalpool.nim | 12 ++++-- 4 files changed, 49 insertions(+), 60 deletions(-) diff --git a/datastore/threadproxyds.nim b/datastore/threadproxyds.nim index d8167f1..d95eb39 100644 --- a/datastore/threadproxyds.nim +++ b/datastore/threadproxyds.nim @@ -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, diff --git a/datastore/threads/threadbackend.nim b/datastore/threads/threadbackend.nim index b6255b7..402d0b5 100644 --- a/datastore/threads/threadbackend.nim +++ b/datastore/threads/threadbackend.nim @@ -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*( diff --git a/datastore/threads/threadresults.nim b/datastore/threads/threadresults.nim index 496b68d..f24fa43 100644 --- a/datastore/threads/threadresults.nim +++ b/datastore/threads/threadresults.nim @@ -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: [].} = diff --git a/datastore/threads/threadsignalpool.nim b/datastore/threads/threadsignalpool.nim index b3e713a..fca66ba 100644 --- a/datastore/threads/threadsignalpool.nim +++ b/datastore/threads/threadsignalpool.nim @@ -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()