mirror of
https://github.com/logos-storage/nim-datastore.git
synced 2026-08-02 12:53:29 +00:00
separate SignalObj -- why isn't it freeing\?
This commit is contained in:
parent
97feed5941
commit
18b0d33ac7
@ -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
|
||||
|
||||
@ -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
|
||||
|
||||
@ -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)
|
||||
|
||||
@ -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
|
||||
|
||||
@ -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()
|
||||
|
||||
Loading…
x
Reference in New Issue
Block a user