From 18b0d33ac70f013a1da538f163f6100f20d48117 Mon Sep 17 00:00:00 2001 From: Jaremy Creechley Date: Tue, 12 Sep 2023 14:51:20 -0700 Subject: [PATCH] separate SignalObj -- why isn't it freeing\? --- datastore/threadproxyds.nim | 41 +++++++++++++------------- datastore/threads/databuffer.nim | 5 ++-- datastore/threads/threadbackend.nim | 31 ++++++++++++------- datastore/threads/threadresults.nim | 9 +++--- datastore/threads/threadsignalpool.nim | 20 +++++++++++++ 5 files changed, 69 insertions(+), 37 deletions(-) diff --git a/datastore/threadproxyds.nim b/datastore/threadproxyds.nim index 0f488b9..75444a9 100644 --- a/datastore/threadproxyds.nim +++ b/datastore/threadproxyds.nim @@ -13,6 +13,7 @@ import ./key import ./query import ./datastore import ./threads/threadbackend +import ./threads/threadsignalpool export key, query @@ -27,13 +28,13 @@ method has*( key: Key ): Future[?!bool] {.async.} = - let ret = await newThreadResult(bool) + let (ret, sig) = await newThreadResult(bool) try: - has(ret, self.tds, key) - await wait(ret[].signal) + has(ret, sig, self.tds, key) + await wait(sig) finally: - ret.release() + discard # ret.release()() return ret.convert(bool) @@ -42,13 +43,13 @@ method delete*( key: Key ): Future[?!void] {.async.} = - let ret = await newThreadResult(void) + let (ret, sig) = await newThreadResult(void) try: - delete(ret, self.tds, key) - await wait(ret[].signal) + delete(ret, sig, self.tds, key) + await wait(sig) finally: - ret.release() + discard # ret.release()() return ret.convert(void) @@ -73,13 +74,13 @@ method get*( ## probably be switched to use a single ThreadSignal ## for the entire batch - let ret = await newThreadResult(ValueBuffer) + let (ret, sig) = await newThreadResult(ValueBuffer) try: - get(ret, self.tds, key) - await wait(ret[].signal) + get(ret, sig, self.tds, key) + await wait(sig) finally: - ret.release() + discard # ret.release()() return ret.convert(seq[byte]) @@ -89,13 +90,13 @@ method put*( data: seq[byte] ): Future[?!void] {.async.} = - let ret = await newThreadResult(void) + let (ret, sig) = await newThreadResult(void) try: - put(ret, self.tds, key, data) - await wait(ret[].signal) + put(ret, sig, self.tds, key, data) + await wait(sig) finally: - ret.release() + discard # ret.release()() return ret.convert(void) @@ -122,7 +123,7 @@ method query*( query: Query ): Future[?!QueryIter] {.async.} = - let ret = await newThreadResult(QueryResponseBuffer) + let (ret, sig) = await newThreadResult(QueryResponseBuffer) # echo "\n\n=== Query Start === " @@ -143,8 +144,8 @@ method query*( iterWrapper.finished = iter[].it.finished if not iter[].it.finished: iterWrapper.readyForNext = false - query(ret, self.tds, iter) - await wait(ret[].signal) + query(ret, sig, self.tds, iter) + await wait(sig) iterWrapper.readyForNext = true # echo "" # print "query:post: ", ret[].results @@ -157,7 +158,7 @@ method query*( proc dispose(): Future[?!void] {.async.} = iter[].it = nil # ensure our sharedptr doesn't try and dealloc - ret.release() + discard # ret.release()() return success() iterWrapper.next = next diff --git a/datastore/threads/databuffer.nim b/datastore/threads/databuffer.nim index 3bfe15d..302c8cb 100644 --- a/datastore/threads/databuffer.nim +++ b/datastore/threads/databuffer.nim @@ -20,12 +20,11 @@ type CatchableErrorBuffer* = object msg: StringBuffer - proc `=destroy`*(x: var DataBufferHolder) = ## copy pointer implementation if x.buf != nil: - # when isMainModule or true: - # echo "buffer: FREE: ", repr x.buf.pointer + when isMainModule or true: + echo "buffer: FREE: ", repr x.buf.pointer deallocShared(x.buf) proc len*(a: DataBuffer): int = a[].size diff --git a/datastore/threads/threadbackend.nim b/datastore/threads/threadbackend.nim index c708950..8fdcd88 100644 --- a/datastore/threads/threadbackend.nim +++ b/datastore/threads/threadbackend.nim @@ -12,6 +12,7 @@ import ../query import ./datastore import ./databuffer import ./threadresults +import ./threadsignalpool # import pretty @@ -70,6 +71,7 @@ type proc hasTask*( ret: TResult[bool], + signal: SharedSignalPtr, tds: ThreadDatastorePtr, kb: KeyBuffer, ) = @@ -82,20 +84,22 @@ proc hasTask*( ret.failure(res.error()) else: ret.success(res.get()) - discard ret[].signal.fireSync() + discard signal.fireSync() except CatchableError as err: ret.failure(err) proc has*( ret: TResult[bool], + signal: SharedSignalPtr, tds: ThreadDatastorePtr, key: Key, ) = let bkey = StringBuffer.new(key.id()) - tds[].tp.spawn hasTask(ret, tds, bkey) + tds[].tp.spawn hasTask(ret, signal, tds, bkey) proc getTask*( ret: TResult[DataBuffer], + signal: SharedSignalPtr, tds: ThreadDatastorePtr, kb: KeyBuffer, ) = @@ -109,21 +113,23 @@ proc getTask*( let db = DataBuffer.new res.get() ret.success(db) - discard ret[].signal.fireSync() + discard signal.fireSync() except CatchableError as err: ret.failure(err) proc get*( ret: TResult[DataBuffer], + signal: SharedSignalPtr, tds: ThreadDatastorePtr, key: Key, ) = let bkey = StringBuffer.new(key.id()) - tds[].tp.spawn getTask(ret, tds, bkey) + tds[].tp.spawn getTask(ret, signal, tds, bkey) proc putTask*( ret: TResult[void], + signal: SharedSignalPtr, tds: ThreadDatastorePtr, kb: KeyBuffer, db: DataBuffer, @@ -140,10 +146,11 @@ proc putTask*( else: ret.success() - discard ret[].signal.fireSync() + discard signal.fireSync() proc put*( ret: TResult[void], + signal: SharedSignalPtr, tds: ThreadDatastorePtr, key: Key, data: seq[byte] @@ -151,11 +158,12 @@ proc put*( let bkey = StringBuffer.new(key.id()) let bval = DataBuffer.new(data) - tds[].tp.spawn putTask(ret, tds, bkey, bval) + tds[].tp.spawn putTask(ret, signal, tds, bkey, bval) proc deleteTask*( ret: TResult[void], + signal: SharedSignalPtr, tds: ThreadDatastorePtr, kb: KeyBuffer, ) = @@ -170,22 +178,24 @@ proc deleteTask*( else: ret.success() - discard ret[].signal.fireSync() + discard signal.fireSync() # import pretty proc delete*( ret: TResult[void], + signal: SharedSignalPtr, tds: ThreadDatastorePtr, key: Key, ) = let bkey = StringBuffer.new(key.id()) - tds[].tp.spawn deleteTask(ret, tds, bkey) + tds[].tp.spawn deleteTask(ret, signal, tds, bkey) # import os proc queryTask*( ret: TResult[QueryResponseBuffer], + signal: SharedSignalPtr, tds: ThreadDatastorePtr, qiter: QueryIterPtr, ) = @@ -205,11 +215,12 @@ proc queryTask*( except Exception as exc: ret.failure(exc) - discard ret[].signal.fireSync() + discard signal.fireSync() proc query*( ret: TResult[QueryResponseBuffer], + signal: SharedSignalPtr, tds: ThreadDatastorePtr, qiter: QueryIterPtr, ) = - tds[].tp.spawn queryTask(ret, tds, qiter) + tds[].tp.spawn queryTask(ret, signal, tds, qiter) diff --git a/datastore/threads/threadresults.nim b/datastore/threads/threadresults.nim index 376ee09..2d8c7f9 100644 --- a/datastore/threads/threadresults.nim +++ b/datastore/threads/threadresults.nim @@ -11,6 +11,7 @@ import ./threadsignalpool export databuffer export smartptrs export threadsync +export threadsignalpool type ThreadSafeTypes* = DataBuffer | void | bool | SharedPtr ##\ @@ -22,7 +23,7 @@ type ## Encapsulates both the results from a thread but also the cross ## thread signaling mechanism. This makes it easier to keep them ## together. - signal*: ThreadSignalPtr + # signal*: SharedSignalPtr results*: Result[T, CatchableErrorBuffer] TResult*[T] = SharedPtr[ThreadResult[T]] ##\ @@ -57,7 +58,7 @@ proc threadSafeType*[T: ThreadSafeTypes](tp: typedesc[T]) = proc newThreadResult*[T]( tp: typedesc[T] -): Future[TResult[T]] {.async.} = +): Future[(TResult[T], SharedSignalPtr)] {.async.} = ## Creates a new TResult including getting ## a new ThreadSignalPtr from the pool. ## @@ -66,8 +67,8 @@ proc newThreadResult*[T]( {.error: "only thread safe types can be used".} let res = newSharedPtr(ThreadResult[T]) - res[].signal = await getThreadSignal() - res + let signal = await newSharedSignalPtr() + (res, signal) proc release*[T](res: TResult[T]) {.raises: [].} = ## release TResult and it's ThreadSignal diff --git a/datastore/threads/threadsignalpool.nim b/datastore/threads/threadsignalpool.nim index d10aefd..4b1da4e 100644 --- a/datastore/threads/threadsignalpool.nim +++ b/datastore/threads/threadsignalpool.nim @@ -76,3 +76,23 @@ proc release*(sig: ThreadSignalPtr) {.raises: [].} = signalPoolUsed.excl(sig) signalPoolFree.incl(sig) # echo "free:signalPoolUsed:size: ", signalPoolUsed.len() + +type + SignalObj* = object + val*: ThreadSignalPtr + + SharedSignalPtr* = SharedPtr[SignalObj] ##\ + +proc `=destroy`*(sig: var SignalObj) = + echo "FREE SIG! ", sig.val.pointer.repr + sig.val.release() + +proc newSharedSignalPtr*(): Future[SharedSignalPtr] {.async, raises: [].} = + let ts = await getThreadSignal() + return newSharedPtr(SignalObj(val: ts)) + +proc fireSync*(sig: SharedSignalPtr): Result[bool, string] = + sig[].val.fireSync() + +proc wait*(sig: SharedSignalPtr): Future[void] = + sig[].val.wait()