renamed CondVar to Semaphore
This commit is contained in:
parent
9df3d9c4d3
commit
81353b2dbc
2 changed files with 27 additions and 27 deletions
|
|
@ -556,7 +556,7 @@ proc wrapProcForSpawn*(owner: PSym; spawnExpr: PNode; retType: PType;
|
||||||
# create flowVar:
|
# create flowVar:
|
||||||
result.add newFastAsgnStmt(fvField, callProc(spawnExpr[2]))
|
result.add newFastAsgnStmt(fvField, callProc(spawnExpr[2]))
|
||||||
if barrier == nil:
|
if barrier == nil:
|
||||||
result.add callCodegenProc("nimFlowVarCreateCondVar", fvField)
|
result.add callCodegenProc("nimFlowVarCreateSemaphore", fvField)
|
||||||
|
|
||||||
elif spawnKind == srByVar:
|
elif spawnKind == srByVar:
|
||||||
var field = newSym(skField, getIdent"fv", owner, n.info)
|
var field = newSym(skField, getIdent"fv", owner, n.info)
|
||||||
|
|
|
||||||
|
|
@ -17,27 +17,27 @@ import cpuinfo, cpuload, locks
|
||||||
{.push stackTrace:off.}
|
{.push stackTrace:off.}
|
||||||
|
|
||||||
type
|
type
|
||||||
CondVar = object
|
Semaphore = object
|
||||||
c: TCond
|
c: TCond
|
||||||
L: TLock
|
L: TLock
|
||||||
counter: int
|
counter: int
|
||||||
|
|
||||||
proc createCondVar(): CondVar =
|
proc createSemaphore(): Semaphore =
|
||||||
initCond(result.c)
|
initCond(result.c)
|
||||||
initLock(result.L)
|
initLock(result.L)
|
||||||
|
|
||||||
proc destroyCondVar(cv: var CondVar) {.inline.} =
|
proc destroySemaphore(cv: var Semaphore) {.inline.} =
|
||||||
deinitCond(cv.c)
|
deinitCond(cv.c)
|
||||||
deinitLock(cv.L)
|
deinitLock(cv.L)
|
||||||
|
|
||||||
proc await(cv: var CondVar) =
|
proc await(cv: var Semaphore) =
|
||||||
acquire(cv.L)
|
acquire(cv.L)
|
||||||
while cv.counter <= 0:
|
while cv.counter <= 0:
|
||||||
wait(cv.c, cv.L)
|
wait(cv.c, cv.L)
|
||||||
dec cv.counter
|
dec cv.counter
|
||||||
release(cv.L)
|
release(cv.L)
|
||||||
|
|
||||||
proc signal(cv: var CondVar) =
|
proc signal(cv: var Semaphore) =
|
||||||
acquire(cv.L)
|
acquire(cv.L)
|
||||||
inc cv.counter
|
inc cv.counter
|
||||||
release(cv.L)
|
release(cv.L)
|
||||||
|
|
@ -48,7 +48,7 @@ const CacheLineSize = 32 # true for most archs
|
||||||
type
|
type
|
||||||
Barrier {.compilerProc.} = object
|
Barrier {.compilerProc.} = object
|
||||||
entered: int
|
entered: int
|
||||||
cv: CondVar # condvar takes 3 words at least
|
cv: Semaphore # Semaphore takes 3 words at least
|
||||||
when sizeof(int) < 8:
|
when sizeof(int) < 8:
|
||||||
cacheAlign: array[CacheLineSize-4*sizeof(int), byte]
|
cacheAlign: array[CacheLineSize-4*sizeof(int), byte]
|
||||||
left: int
|
left: int
|
||||||
|
|
@ -75,12 +75,12 @@ proc openBarrier(b: ptr Barrier) {.compilerProc, inline.} =
|
||||||
proc closeBarrier(b: ptr Barrier) {.compilerProc.} =
|
proc closeBarrier(b: ptr Barrier) {.compilerProc.} =
|
||||||
fence()
|
fence()
|
||||||
if b.left != b.entered:
|
if b.left != b.entered:
|
||||||
b.cv = createCondVar()
|
b.cv = createSemaphore()
|
||||||
fence()
|
fence()
|
||||||
b.interest = true
|
b.interest = true
|
||||||
fence()
|
fence()
|
||||||
while b.left != b.entered: await(b.cv)
|
while b.left != b.entered: await(b.cv)
|
||||||
destroyCondVar(b.cv)
|
destroySemaphore(b.cv)
|
||||||
|
|
||||||
{.pop.}
|
{.pop.}
|
||||||
|
|
||||||
|
|
@ -90,13 +90,13 @@ type
|
||||||
foreign* = object ## a region that indicates the pointer comes from a
|
foreign* = object ## a region that indicates the pointer comes from a
|
||||||
## foreign thread heap.
|
## foreign thread heap.
|
||||||
AwaitInfo = object
|
AwaitInfo = object
|
||||||
cv: CondVar
|
cv: Semaphore
|
||||||
idx: int
|
idx: int
|
||||||
|
|
||||||
FlowVarBase* = ref FlowVarBaseObj ## untyped base class for 'FlowVar[T]'
|
FlowVarBase* = ref FlowVarBaseObj ## untyped base class for 'FlowVar[T]'
|
||||||
FlowVarBaseObj = object of RootObj
|
FlowVarBaseObj = object of RootObj
|
||||||
ready, usesCondVar, awaited: bool
|
ready, usesSemaphore, awaited: bool
|
||||||
cv: CondVar #\
|
cv: Semaphore #\
|
||||||
# for 'awaitAny' support
|
# for 'awaitAny' support
|
||||||
ai: ptr AwaitInfo
|
ai: ptr AwaitInfo
|
||||||
idx: int
|
idx: int
|
||||||
|
|
@ -112,13 +112,13 @@ type
|
||||||
ToFreeQueue = object
|
ToFreeQueue = object
|
||||||
len: int
|
len: int
|
||||||
lock: TLock
|
lock: TLock
|
||||||
empty: CondVar
|
empty: Semaphore
|
||||||
data: array[128, pointer]
|
data: array[128, pointer]
|
||||||
|
|
||||||
WorkerProc = proc (thread, args: pointer) {.nimcall, gcsafe.}
|
WorkerProc = proc (thread, args: pointer) {.nimcall, gcsafe.}
|
||||||
Worker = object
|
Worker = object
|
||||||
taskArrived: CondVar
|
taskArrived: Semaphore
|
||||||
taskStarted: CondVar #\
|
taskStarted: Semaphore #\
|
||||||
# task data:
|
# task data:
|
||||||
f: WorkerProc
|
f: WorkerProc
|
||||||
data: pointer
|
data: pointer
|
||||||
|
|
@ -130,10 +130,10 @@ type
|
||||||
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
|
||||||
## to call this explicitly.
|
## to call this explicitly.
|
||||||
if fv.usesCondVar and not fv.awaited:
|
if fv.usesSemaphore and not fv.awaited:
|
||||||
fv.awaited = true
|
fv.awaited = true
|
||||||
await(fv.cv)
|
await(fv.cv)
|
||||||
destroyCondVar(fv.cv)
|
destroySemaphore(fv.cv)
|
||||||
|
|
||||||
proc selectWorker(w: ptr Worker; fn: WorkerProc; data: pointer): bool =
|
proc selectWorker(w: ptr Worker; fn: WorkerProc; data: pointer): bool =
|
||||||
if cas(addr w.ready, true, false):
|
if cas(addr w.ready, true, false):
|
||||||
|
|
@ -191,9 +191,9 @@ proc fvFinalizer[T](fv: FlowVar[T]) = finished(fv)
|
||||||
proc nimCreateFlowVar[T](): FlowVar[T] {.compilerProc.} =
|
proc nimCreateFlowVar[T](): FlowVar[T] {.compilerProc.} =
|
||||||
new(result, fvFinalizer)
|
new(result, fvFinalizer)
|
||||||
|
|
||||||
proc nimFlowVarCreateCondVar(fv: FlowVarBase) {.compilerProc.} =
|
proc nimFlowVarCreateSemaphore(fv: FlowVarBase) {.compilerProc.} =
|
||||||
fv.cv = createCondVar()
|
fv.cv = createSemaphore()
|
||||||
fv.usesCondVar = true
|
fv.usesSemaphore = true
|
||||||
|
|
||||||
proc nimFlowVarSignal(fv: FlowVarBase) {.compilerProc.} =
|
proc nimFlowVarSignal(fv: FlowVarBase) {.compilerProc.} =
|
||||||
if fv.ai != nil:
|
if fv.ai != nil:
|
||||||
|
|
@ -202,7 +202,7 @@ proc nimFlowVarSignal(fv: FlowVarBase) {.compilerProc.} =
|
||||||
inc fv.ai.cv.counter
|
inc fv.ai.cv.counter
|
||||||
release(fv.ai.cv.L)
|
release(fv.ai.cv.L)
|
||||||
signal(fv.ai.cv.c)
|
signal(fv.ai.cv.c)
|
||||||
if fv.usesCondVar:
|
if fv.usesSemaphore:
|
||||||
signal(fv.cv)
|
signal(fv.cv)
|
||||||
|
|
||||||
proc awaitAndThen*[T](fv: FlowVar[T]; action: proc (x: T) {.closure.}) =
|
proc awaitAndThen*[T](fv: FlowVar[T]; action: proc (x: T) {.closure.}) =
|
||||||
|
|
@ -242,7 +242,7 @@ proc awaitAny*(flowVars: openArray[FlowVarBase]): int =
|
||||||
## **Note**: This results in non-deterministic behaviour and so should be
|
## **Note**: This results in non-deterministic behaviour and so should be
|
||||||
## avoided.
|
## avoided.
|
||||||
var ai: AwaitInfo
|
var ai: AwaitInfo
|
||||||
ai.cv = createCondVar()
|
ai.cv = createSemaphore()
|
||||||
var conflicts = 0
|
var conflicts = 0
|
||||||
for i in 0 .. flowVars.high:
|
for i in 0 .. flowVars.high:
|
||||||
if cas(addr flowVars[i].ai, nil, addr ai):
|
if cas(addr flowVars[i].ai, nil, addr ai):
|
||||||
|
|
@ -256,7 +256,7 @@ proc awaitAny*(flowVars: openArray[FlowVarBase]): int =
|
||||||
discard cas(addr flowVars[i].ai, addr ai, nil)
|
discard cas(addr flowVars[i].ai, addr ai, nil)
|
||||||
else:
|
else:
|
||||||
result = -1
|
result = -1
|
||||||
destroyCondVar(ai.cv)
|
destroySemaphore(ai.cv)
|
||||||
|
|
||||||
proc nimArgsPassingDone(p: pointer) {.compilerProc.} =
|
proc nimArgsPassingDone(p: pointer) {.compilerProc.} =
|
||||||
let w = cast[ptr Worker](p)
|
let w = cast[ptr Worker](p)
|
||||||
|
|
@ -270,7 +270,7 @@ var
|
||||||
currentPoolSize: int
|
currentPoolSize: int
|
||||||
maxPoolSize = MaxThreadPoolSize
|
maxPoolSize = MaxThreadPoolSize
|
||||||
minPoolSize = 4
|
minPoolSize = 4
|
||||||
gSomeReady = createCondVar()
|
gSomeReady = createSemaphore()
|
||||||
readyWorker: ptr Worker
|
readyWorker: ptr Worker
|
||||||
|
|
||||||
proc slave(w: ptr Worker) {.thread.} =
|
proc slave(w: ptr Worker) {.thread.} =
|
||||||
|
|
@ -307,10 +307,10 @@ proc setMaxPoolSize*(size: range[1..MaxThreadPoolSize]) =
|
||||||
w.shutdown = true
|
w.shutdown = true
|
||||||
|
|
||||||
proc activateThread(i: int) {.noinline.} =
|
proc activateThread(i: int) {.noinline.} =
|
||||||
workersData[i].taskArrived = createCondVar()
|
workersData[i].taskArrived = createSemaphore()
|
||||||
workersData[i].taskStarted = createCondVar()
|
workersData[i].taskStarted = createSemaphore()
|
||||||
workersData[i].initialized = true
|
workersData[i].initialized = true
|
||||||
workersData[i].q.empty = createCondVar()
|
workersData[i].q.empty = createSemaphore()
|
||||||
initLock(workersData[i].q.lock)
|
initLock(workersData[i].q.lock)
|
||||||
createThread(workers[i], slave, addr(workersData[i]))
|
createThread(workers[i], slave, addr(workersData[i]))
|
||||||
|
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue