fixed the deadlock that happens when stress testing ToFreeQueue
This commit is contained in:
parent
06e9932e8a
commit
943d4ee714
2 changed files with 31 additions and 26 deletions
|
|
@ -112,8 +112,8 @@ type
|
||||||
ToFreeQueue = object
|
ToFreeQueue = object
|
||||||
len: int
|
len: int
|
||||||
lock: TLock
|
lock: TLock
|
||||||
empty: TCond
|
empty: CondVar
|
||||||
data: array[2, pointer]
|
data: array[128, pointer]
|
||||||
|
|
||||||
WorkerProc = proc (thread, args: pointer) {.nimcall, gcsafe.}
|
WorkerProc = proc (thread, args: pointer) {.nimcall, gcsafe.}
|
||||||
Worker = object
|
Worker = object
|
||||||
|
|
@ -143,12 +143,26 @@ proc selectWorker(w: ptr Worker; fn: WorkerProc; data: pointer): bool =
|
||||||
await(w.taskStarted)
|
await(w.taskStarted)
|
||||||
result = true
|
result = true
|
||||||
|
|
||||||
|
proc cleanFlowVars(w: ptr Worker) =
|
||||||
|
let q = addr(w.q)
|
||||||
|
acquire(q.lock)
|
||||||
|
for i in 0 .. <q.len:
|
||||||
|
GC_unref(cast[RootRef](q.data[i]))
|
||||||
|
#echo "GC_unref"
|
||||||
|
q.len = 0
|
||||||
|
release(q.lock)
|
||||||
|
|
||||||
proc wakeupWorkerToProcessQueue(w: ptr Worker) =
|
proc wakeupWorkerToProcessQueue(w: ptr Worker) =
|
||||||
# Note that if this fails somebody else already woke up the thread so it's
|
# we have to ensure it's us who wakes up the owning thread.
|
||||||
# perfectly fine to do nothing:
|
# This is quite horrible code, but it runs so rarely that it doesn't matter:
|
||||||
if cas(addr w.ready, true, false):
|
while not cas(addr w.ready, true, false):
|
||||||
|
cpuRelax()
|
||||||
|
discard
|
||||||
w.data = nil
|
w.data = nil
|
||||||
w.f = proc (t, a: pointer) {.nimcall.} = discard
|
w.f = proc (w, a: pointer) {.nimcall.} =
|
||||||
|
let w = cast[ptr Worker](w)
|
||||||
|
cleanFlowVars(w)
|
||||||
|
signal(w.q.empty)
|
||||||
signal(w.taskArrived)
|
signal(w.taskArrived)
|
||||||
|
|
||||||
proc finished(fv: FlowVarBase) =
|
proc finished(fv: FlowVarBase) =
|
||||||
|
|
@ -160,29 +174,17 @@ proc finished(fv: FlowVarBase) =
|
||||||
if fv.data.isNil: return
|
if fv.data.isNil: return
|
||||||
let owner = cast[ptr Worker](fv.owner)
|
let owner = cast[ptr Worker](fv.owner)
|
||||||
let q = addr(owner.q)
|
let q = addr(owner.q)
|
||||||
var waited = false
|
|
||||||
acquire(q.lock)
|
acquire(q.lock)
|
||||||
while not (q.len < q.data.len):
|
while not (q.len < q.data.len):
|
||||||
#echo "EXHAUSTED!"
|
#echo "EXHAUSTED!"
|
||||||
|
release(q.lock)
|
||||||
wakeupWorkerToProcessQueue(owner)
|
wakeupWorkerToProcessQueue(owner)
|
||||||
wait(q.empty, q.lock)
|
await(q.empty)
|
||||||
waited = true
|
acquire(q.lock)
|
||||||
q.data[q.len] = cast[pointer](fv.data)
|
q.data[q.len] = cast[pointer](fv.data)
|
||||||
inc q.len
|
inc q.len
|
||||||
release(q.lock)
|
release(q.lock)
|
||||||
fv.data = nil
|
fv.data = nil
|
||||||
# wakeup other potentially waiting threads:
|
|
||||||
if waited: signal(q.empty)
|
|
||||||
|
|
||||||
proc cleanFlowVars(w: ptr Worker) =
|
|
||||||
let q = addr(w.q)
|
|
||||||
acquire(q.lock)
|
|
||||||
for i in 0 .. <q.len:
|
|
||||||
GC_unref(cast[RootRef](q.data[i]))
|
|
||||||
#echo "GC_unref"
|
|
||||||
q.len = 0
|
|
||||||
release(q.lock)
|
|
||||||
signal(q.empty)
|
|
||||||
|
|
||||||
proc fvFinalizer[T](fv: FlowVar[T]) = finished(fv)
|
proc fvFinalizer[T](fv: FlowVar[T]) = finished(fv)
|
||||||
|
|
||||||
|
|
@ -273,6 +275,9 @@ var
|
||||||
|
|
||||||
proc slave(w: ptr Worker) {.thread.} =
|
proc slave(w: ptr Worker) {.thread.} =
|
||||||
while true:
|
while true:
|
||||||
|
when declared(atomicStoreN):
|
||||||
|
atomicStoreN(addr(w.ready), true, ATOMIC_SEQ_CST)
|
||||||
|
else:
|
||||||
w.ready = true
|
w.ready = true
|
||||||
readyWorker = w
|
readyWorker = w
|
||||||
signal(gSomeReady)
|
signal(gSomeReady)
|
||||||
|
|
@ -305,7 +310,7 @@ proc activateThread(i: int) {.noinline.} =
|
||||||
workersData[i].taskArrived = createCondVar()
|
workersData[i].taskArrived = createCondVar()
|
||||||
workersData[i].taskStarted = createCondVar()
|
workersData[i].taskStarted = createCondVar()
|
||||||
workersData[i].initialized = true
|
workersData[i].initialized = true
|
||||||
initCond(workersData[i].q.empty)
|
workersData[i].q.empty = createCondVar()
|
||||||
initLock(workersData[i].q.lock)
|
initLock(workersData[i].q.lock)
|
||||||
createThread(workers[i], slave, addr(workersData[i]))
|
createThread(workers[i], slave, addr(workersData[i]))
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -1,7 +1,7 @@
|
||||||
discard """
|
discard """
|
||||||
output: '''called deepCopy for int
|
output: '''called deepCopy for int
|
||||||
called deepCopy for int
|
called deepCopy for int
|
||||||
done999 999
|
done999 999'''
|
||||||
"""
|
"""
|
||||||
|
|
||||||
import threadpool
|
import threadpool
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue