This commit is contained in:
Jaremy Creechley 2023-09-12 17:22:04 -07:00
parent 08c4a936ad
commit cf891184a9
No known key found for this signature in database
GPG Key ID: 4E66FB67B21D3300
5 changed files with 71 additions and 44 deletions

View File

@ -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,

View File

@ -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)

View File

@ -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,

View File

@ -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..<x.len: `=destroy`(x.data[i])
echo "SIGNAL: FREE: ", repr x.buf.pointer, " ", x.cnt[]
deallocShared(x.buf)
deallocShared(x.cnt)
else:
echo "SIGNAL: decr: ", repr x.buf.pointer, " ", x.cnt[]
proc `=destroy`*(x: var SharedSignalPtr) =
echo "SIGNAL: destroy: ", repr x.buf.pointer, " ", x.cnt.repr
echo "SIGNAL: destroy:st: ", $getStackTrace()
decr(x)
proc `=copy`*(a: var SharedSignalPtr; b: SharedSignalPtr) =
# do nothing for self-assignments:
if a.buf == b.buf: return
@ -108,16 +116,11 @@ proc `=copy`*(a: var SharedSignalPtr; b: SharedSignalPtr) =
discard atomicAddFetch(b.cnt, 1, ATOMIC_RELAXED)
a.buf = b.buf
a.cnt = b.cnt
echo "SIGNAL: Copy: repr: ", b.cnt[],
" ", repr a.buf.pointer,
" ", repr b.buf.pointer
proc `incr`*(a: SharedSignalPtr) =
echo "SIGNAL: incr: ", atomicAddFetch(a.cnt, 1, ATOMIC_RELAXED)
proc newSharedSignalPtr*(): Future[SharedSignalPtr] {.async, raises: [].} =
result.cnt = cast[ptr int](allocShared0(sizeof(result.cnt)))
result.buf = await getThreadSignal()
echo "SIGNAL: alloc: ", repr result.buf.pointer
template fireSync*(sig: SharedSignalPtr): untyped =
let ts: ThreadSignalPtr = sig.buf

View File

@ -16,34 +16,40 @@ import ./querycommontests
# import pretty
suite "Test Basic ThreadProxyDatastore":
var
sds: ThreadProxyDatastore
mem: MemoryDatastore
key1: Key
data: seq[byte]
proc testThreadProxy() =
setupAll:
mem = MemoryDatastore.new()
sds = newThreadProxyDatastore(mem).expect("should work")
key1 = Key.init("/a").tryGet
data = "value for 1".toBytes()
suite "Test Basic ThreadProxyDatastore":
var
sds: ThreadProxyDatastore
mem: MemoryDatastore
key1: Key
data: seq[byte]
test "check put":
# echo "\n\n=== put ==="
let res1 = await sds.put(key1, data)
check res1.isOk
# print "res1: ", res1
setupAll:
mem = MemoryDatastore.new()
sds = newThreadProxyDatastore(mem).expect("should work")
key1 = Key.init("/a").tryGet
data = "value for 1".toBytes()
# test "check get":
# # echo "\n\n=== get ==="
# let res2 = await sds.get(key1)
# check res2.get() == data
# var val = ""
# for c in res2.get():
# val &= char(c)
# # print "get res2: ", $val
test "check put":
for i in 1..2:
echo "\n\n=== put ==="
let res1 = await sds.put(key1, data)
check res1.isOk
# GC_fullCollect()
# print "res1: ", res1
# test "check get":
# # echo "\n\n=== get ==="
# let res2 = await sds.get(key1)
# check res2.get() == data
# var val = ""
# for c in res2.get():
# val &= char(c)
# # print "get res2: ", $val
testThreadProxy()
GC_fullCollect()
# suite "Test Basic ThreadProxyDatastore":