setMaxPoolSize improvements
This commit is contained in:
parent
057b53e139
commit
68accb53c4
1 changed files with 10 additions and 5 deletions
|
|
@ -134,7 +134,7 @@ proc await*(fv: FlowVarBase) =
|
||||||
proc finished(fv: FlowVarBase) =
|
proc finished(fv: FlowVarBase) =
|
||||||
doAssert fv.ai.isNil, "flowVar is still attached to an 'awaitAny'"
|
doAssert fv.ai.isNil, "flowVar is still attached to an 'awaitAny'"
|
||||||
# we have to protect against the rare cases where the owner of the flowVar
|
# we have to protect against the rare cases where the owner of the flowVar
|
||||||
# simply disregards the flowVar and yet the "flowVarr" has not yet written
|
# simply disregards the flowVar and yet the "flowVar" has not yet written
|
||||||
# anything to it:
|
# anything to it:
|
||||||
await(fv)
|
await(fv)
|
||||||
if fv.data.isNil: return
|
if fv.data.isNil: return
|
||||||
|
|
@ -207,6 +207,7 @@ proc `^`*[T](fv: FlowVar[T]): T =
|
||||||
## blocks until the value is available and then returns this value.
|
## blocks until the value is available and then returns this value.
|
||||||
await(fv)
|
await(fv)
|
||||||
when T is string or T is seq:
|
when T is string or T is seq:
|
||||||
|
# XXX closures? deepCopy?
|
||||||
result = cast[T](fv.data)
|
result = cast[T](fv.data)
|
||||||
else:
|
else:
|
||||||
result = fv.blob
|
result = fv.blob
|
||||||
|
|
@ -264,6 +265,10 @@ proc slave(w: ptr Worker) {.thread.} =
|
||||||
w.shutdown = false
|
w.shutdown = false
|
||||||
atomicDec currentPoolSize
|
atomicDec currentPoolSize
|
||||||
|
|
||||||
|
var
|
||||||
|
workers: array[MaxThreadPoolSize, TThread[ptr Worker]]
|
||||||
|
workersData: array[MaxThreadPoolSize, Worker]
|
||||||
|
|
||||||
proc setMinPoolSize*(size: range[1..MaxThreadPoolSize]) =
|
proc setMinPoolSize*(size: range[1..MaxThreadPoolSize]) =
|
||||||
## sets the minimal thread pool size. The default value of this is 4.
|
## sets the minimal thread pool size. The default value of this is 4.
|
||||||
minPoolSize = size
|
minPoolSize = size
|
||||||
|
|
@ -272,10 +277,10 @@ proc setMaxPoolSize*(size: range[1..MaxThreadPoolSize]) =
|
||||||
## sets the minimal thread pool size. The default value of this
|
## sets the minimal thread pool size. The default value of this
|
||||||
## is ``MaxThreadPoolSize``.
|
## is ``MaxThreadPoolSize``.
|
||||||
maxPoolSize = size
|
maxPoolSize = size
|
||||||
|
if currentPoolSize > maxPoolSize:
|
||||||
var
|
for i in maxPoolSize..currentPoolSize-1:
|
||||||
workers: array[MaxThreadPoolSize, TThread[ptr Worker]]
|
let w = addr(workersData[i])
|
||||||
workersData: array[MaxThreadPoolSize, Worker]
|
w.shutdown = true
|
||||||
|
|
||||||
proc activateThread(i: int) {.noinline.} =
|
proc activateThread(i: int) {.noinline.} =
|
||||||
workersData[i].taskArrived = createCondVar()
|
workersData[i].taskArrived = createCondVar()
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue