proper waiting for the pinned thread
This commit is contained in:
parent
21ea8e6913
commit
3312d49a48
1 changed files with 6 additions and 3 deletions
|
|
@ -128,6 +128,7 @@ type
|
||||||
initialized: bool # whether it has even been initialized
|
initialized: bool # whether it has even been initialized
|
||||||
shutdown: bool # the pool requests to shut down this worker thread
|
shutdown: bool # the pool requests to shut down this worker thread
|
||||||
q: ToFreeQueue
|
q: ToFreeQueue
|
||||||
|
readyForTask: Semaphore
|
||||||
|
|
||||||
proc await*(fv: FlowVarBase) =
|
proc await*(fv: FlowVarBase) =
|
||||||
## waits until the value for the flowVar arrives. Usually it is not necessary
|
## waits until the value for the flowVar arrives. Usually it is not necessary
|
||||||
|
|
@ -301,6 +302,7 @@ proc distinguishedSlave(w: ptr Worker) {.thread.} =
|
||||||
atomicStoreN(addr(w.ready), true, ATOMIC_SEQ_CST)
|
atomicStoreN(addr(w.ready), true, ATOMIC_SEQ_CST)
|
||||||
else:
|
else:
|
||||||
w.ready = true
|
w.ready = true
|
||||||
|
signal(w.readyForTask)
|
||||||
await(w.taskArrived)
|
await(w.taskArrived)
|
||||||
assert(not w.ready)
|
assert(not w.ready)
|
||||||
w.f(w, w.data)
|
w.f(w, w.data)
|
||||||
|
|
@ -340,6 +342,7 @@ proc activateDistinguishedThread(i: int) {.noinline.} =
|
||||||
distinguishedData[i].initialized = true
|
distinguishedData[i].initialized = true
|
||||||
distinguishedData[i].q.empty = createSemaphore()
|
distinguishedData[i].q.empty = createSemaphore()
|
||||||
initLock(distinguishedData[i].q.lock)
|
initLock(distinguishedData[i].q.lock)
|
||||||
|
distinguishedData[i].readyForTask = createSemaphore()
|
||||||
createThread(distinguished[i], distinguishedSlave, addr(distinguishedData[i]))
|
createThread(distinguished[i], distinguishedSlave, addr(distinguishedData[i]))
|
||||||
|
|
||||||
proc setup() =
|
proc setup() =
|
||||||
|
|
@ -429,11 +432,11 @@ proc nimSpawn4(fn: WorkerProc; data: pointer; id: ThreadId) {.compilerProc.} =
|
||||||
acquire(distinguishedLock)
|
acquire(distinguishedLock)
|
||||||
if not distinguishedData[id].initialized:
|
if not distinguishedData[id].initialized:
|
||||||
activateDistinguishedThread(id)
|
activateDistinguishedThread(id)
|
||||||
|
release(distinguishedLock)
|
||||||
while true:
|
while true:
|
||||||
if selectWorker(addr(distinguishedData[id]), fn, data): break
|
if selectWorker(addr(distinguishedData[id]), fn, data): break
|
||||||
cpuRelax()
|
await(distinguishedData[id].readyForTask)
|
||||||
# XXX exponential backoff?
|
|
||||||
release(distinguishedLock)
|
|
||||||
|
|
||||||
proc sync*() =
|
proc sync*() =
|
||||||
## a simple barrier to wait for all spawn'ed tasks. If you need more elaborate
|
## a simple barrier to wait for all spawn'ed tasks. If you need more elaborate
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue