diff --git a/datastore/threadproxyds.nim b/datastore/threadproxyds.nim index 75444a9..b9d6580 100644 --- a/datastore/threadproxyds.nim +++ b/datastore/threadproxyds.nim @@ -32,7 +32,7 @@ method has*( try: has(ret, sig, self.tds, key) - await wait(sig) + wait(sig) finally: discard # ret.release()() @@ -47,7 +47,7 @@ method delete*( try: delete(ret, sig, self.tds, key) - await wait(sig) + wait(sig) finally: discard # ret.release()() @@ -78,7 +78,7 @@ method get*( try: get(ret, sig, self.tds, key) - await wait(sig) + wait(sig) finally: discard # ret.release()() @@ -90,13 +90,16 @@ method put*( data: seq[byte] ): Future[?!void] {.async.} = - let (ret, sig) = await newThreadResult(void) + var ret: TResult[void] + block: + let (rets, sig) = await newThreadResult(void) + ret = rets - try: - put(ret, sig, self.tds, key, data) - await wait(sig) - finally: - discard # ret.release()() + try: + put(ret, sig, self.tds, key, data) + wait(sig) + finally: + discard # ret.release()() return ret.convert(void) @@ -145,7 +148,7 @@ method query*( if not iter[].it.finished: iterWrapper.readyForNext = false query(ret, sig, self.tds, iter) - await wait(sig) + wait(sig) iterWrapper.readyForNext = true # echo "" # print "query:post: ", ret[].results diff --git a/datastore/threads/threadresults.nim b/datastore/threads/threadresults.nim index 2d8c7f9..3b0b219 100644 --- a/datastore/threads/threadresults.nim +++ b/datastore/threads/threadresults.nim @@ -30,15 +30,6 @@ type ## SharedPtr that allocates a shared buffer and keeps the ## memory allocated until all references to it are gone. ## - ## Important: - ## On `refc` that "internal" destructors for ThreadResult[T] - ## are *not* called. Effectively limiting this to 1 depth - ## of destructors. Hence the `threadSafeType` marker below. - ## - ## Edit: not sure this is quire accurate, but some care - ## needs to be taken to verify the destructor - ## works with the specific type. - ## ## Since ThreadResult is a plain object, its lifetime can be ## tied to that of an async proc. In this case it could be ## freed before the other background thread is finished. diff --git a/datastore/threads/threadsignalpool.nim b/datastore/threads/threadsignalpool.nim index 4b1da4e..fa13d8d 100644 --- a/datastore/threads/threadsignalpool.nim +++ b/datastore/threads/threadsignalpool.nim @@ -77,22 +77,52 @@ proc release*(sig: ThreadSignalPtr) {.raises: [].} = signalPoolFree.incl(sig) # echo "free:signalPoolUsed:size: ", signalPoolUsed.len() + type - SignalObj* = object - val*: ThreadSignalPtr + SharedSignalPtr* = object + cnt: ptr int + buf*: ThreadSignalPtr - SharedSignalPtr* = SharedPtr[SignalObj] ##\ -proc `=destroy`*(sig: var SignalObj) = - echo "FREE SIG! ", sig.val.pointer.repr - sig.val.release() +proc `$`*(data: SharedSignalPtr): string = + if data.buf.isNil: + result = "nil" + else: + result = data.buf.pointer.repr + +proc `=destroy`*(x: var SharedSignalPtr) = + if x.buf != nil and x.cnt != nil: + let res = atomicSubFetch(x.cnt, 1, ATOMIC_ACQUIRE) + if res == 0: + # for i in 0.. 0 +# check res.len() > 0