diff --git a/datastore/threadproxyds.nim b/datastore/threadproxyds.nim index b9d6580..ffd8b02 100644 --- a/datastore/threadproxyds.nim +++ b/datastore/threadproxyds.nim @@ -84,24 +84,36 @@ method get*( return ret.convert(seq[byte]) +type + FutureAlloc* = ref object + data*: int + +proc faFinalizer*(fa: FutureAlloc) = + echo "FA FREE! ", cast[pointer](fa).repr + method put*( self: ThreadProxyDatastore, key: Key, data: seq[byte] ): Future[?!void] {.async.} = - var ret: TResult[void] block: - let (rets, sig) = await newThreadResult(void) - ret = rets + var fa: FutureAlloc + fa.new(faFinalizer) + + echo "FA NEW! ", cast[pointer](fa).repr + + let (ret, reallysignal) = await newThreadResult(void) try: - put(ret, sig, self.tds, key, data) - wait(sig) + puts(ret, reallysignal, self.tds, key, data) + wait(reallysignal) finally: + reallysignal.decr() discard # ret.release()() - return ret.convert(void) + result = ret.convert(void) + method put*( self: ThreadProxyDatastore, diff --git a/datastore/threads/databuffer.nim b/datastore/threads/databuffer.nim index 302c8cb..86bee22 100644 --- a/datastore/threads/databuffer.nim +++ b/datastore/threads/databuffer.nim @@ -41,17 +41,20 @@ proc `==`*(a, b: DataBuffer): bool = elif a[].buf == b[].buf: return true else: a.hash() == b.hash() -proc new*(tp: typedesc[DataBuffer], size: int = 0): DataBuffer = +proc new*[T: DataBuffer](tp: typedesc[T], size: int = 0): T = ## allocate new buffer with given size - newSharedPtr(DataBufferHolder( + result = newSharedPtr(DataBufferHolder( buf: cast[typeof(result[].buf)](allocShared0(size)), size: size, )) + echo "DataBuffer: ALLOC: ", repr result[].buf.pointer, + " kind: ", $(typeof(result)), + " sharePtr: ", result.val.pointer.repr -proc new*[T: byte | char](tp: typedesc[DataBuffer], data: openArray[T]): DataBuffer = +proc new*[T: byte | char; B](tp: typedesc[B], data: openArray[T]): DataBuffer = ## allocate new buffer and copies indata from openArray ## - result = DataBuffer.new(data.len) + result = B.new(data.len) if data.len() > 0: copyMem(result[].buf, unsafeAddr data[0], data.len) diff --git a/datastore/threads/threadbackend.nim b/datastore/threads/threadbackend.nim index 8fdcd88..fe719be 100644 --- a/datastore/threads/threadbackend.nim +++ b/datastore/threads/threadbackend.nim @@ -147,8 +147,11 @@ proc putTask*( ret.success() discard signal.fireSync() + echo "PUT DONE:kb: ", kb.val.pointer.repr + echo "PUT DONE:db: ", db.val.pointer.repr + # GC_fullCollect() -proc put*( +proc puts*( ret: TResult[void], signal: SharedSignalPtr, tds: ThreadDatastorePtr, diff --git a/datastore/threads/threadsignalpool.nim b/datastore/threads/threadsignalpool.nim index fa13d8d..158fd23 100644 --- a/datastore/threads/threadsignalpool.nim +++ b/datastore/threads/threadsignalpool.nim @@ -90,17 +90,25 @@ proc `$`*(data: SharedSignalPtr): string = else: result = data.buf.pointer.repr -proc `=destroy`*(x: var SharedSignalPtr) = +proc `incr`*(a: SharedSignalPtr) = + echo "SIGNAL: manual incr: ", atomicAddFetch(a.cnt, 1, ATOMIC_RELAXED) + +proc `decr`*(x: SharedSignalPtr) = if x.buf != nil and x.cnt != nil: let res = atomicSubFetch(x.cnt, 1, ATOMIC_ACQUIRE) if res == 0: - # for i in 0..