fixes #7638; awaitAny blocks if the flow vars all have been complete already
This commit is contained in:
parent
cb03ae2c9f
commit
0dc4d6dcc2
2 changed files with 52 additions and 7 deletions
|
|
@ -168,6 +168,15 @@ proc wakeupWorkerToProcessQueue(w: ptr Worker) =
|
||||||
signal(w.q.empty)
|
signal(w.q.empty)
|
||||||
signal(w.taskArrived)
|
signal(w.taskArrived)
|
||||||
|
|
||||||
|
proc attach(fv: FlowVarBase; i: int): bool =
|
||||||
|
acquire(fv.cv.L)
|
||||||
|
if fv.cv.counter <= 0:
|
||||||
|
fv.idx = i
|
||||||
|
result = true
|
||||||
|
else:
|
||||||
|
result = false
|
||||||
|
release(fv.cv.L)
|
||||||
|
|
||||||
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
|
||||||
|
|
@ -248,23 +257,24 @@ proc awaitAny*(flowVars: openArray[FlowVarBase]): int =
|
||||||
## the same time. That means if you awaitAny([a,b]) and awaitAny([b,c]) the second
|
## the same time. That means if you awaitAny([a,b]) and awaitAny([b,c]) the second
|
||||||
## call will only await 'c'. If there is no flowVar left to be able to wait
|
## call will only await 'c'. If there is no flowVar left to be able to wait
|
||||||
## on, -1 is returned.
|
## on, -1 is returned.
|
||||||
## **Note**: This results in non-deterministic behaviour and so should be
|
## **Note**: This results in non-deterministic behaviour and should be avoided.
|
||||||
## avoided.
|
|
||||||
var ai: AwaitInfo
|
var ai: AwaitInfo
|
||||||
ai.cv.initSemaphore()
|
ai.cv.initSemaphore()
|
||||||
var conflicts = 0
|
var conflicts = 0
|
||||||
|
result = -1
|
||||||
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):
|
||||||
flowVars[i].idx = i
|
if not attach(flowVars[i], i):
|
||||||
|
result = i
|
||||||
|
break
|
||||||
else:
|
else:
|
||||||
inc conflicts
|
inc conflicts
|
||||||
if conflicts < flowVars.len:
|
if conflicts < flowVars.len:
|
||||||
await(ai.cv)
|
if result < 0:
|
||||||
result = ai.idx
|
await(ai.cv)
|
||||||
|
result = ai.idx
|
||||||
for i in 0 .. flowVars.high:
|
for i in 0 .. flowVars.high:
|
||||||
discard cas(addr flowVars[i].ai, addr ai, nil)
|
discard cas(addr flowVars[i].ai, addr ai, nil)
|
||||||
else:
|
|
||||||
result = -1
|
|
||||||
destroySemaphore(ai.cv)
|
destroySemaphore(ai.cv)
|
||||||
|
|
||||||
proc isReady*(fv: FlowVarBase): bool =
|
proc isReady*(fv: FlowVarBase): bool =
|
||||||
|
|
|
||||||
35
tests/parallel/twaitany.nim
Normal file
35
tests/parallel/twaitany.nim
Normal file
|
|
@ -0,0 +1,35 @@
|
||||||
|
discard """
|
||||||
|
output: '''true'''
|
||||||
|
"""
|
||||||
|
|
||||||
|
# bug #7638
|
||||||
|
import threadpool, os, strformat
|
||||||
|
|
||||||
|
proc timer(d: int): int =
|
||||||
|
#echo fmt"sleeping {d}"
|
||||||
|
sleep(d)
|
||||||
|
#echo fmt"done {d}"
|
||||||
|
return d
|
||||||
|
|
||||||
|
var durations = [1000, 2000, 3000, 4000, 5000]
|
||||||
|
var tasks: seq[FlowVarBase] = @[]
|
||||||
|
var results: seq[int] = @[]
|
||||||
|
|
||||||
|
for i in 0 .. durations.high:
|
||||||
|
tasks.add spawn timer(durations[i])
|
||||||
|
|
||||||
|
var index = awaitAny(tasks)
|
||||||
|
while index != -1:
|
||||||
|
results.add ^cast[FlowVar[int]](tasks[index])
|
||||||
|
tasks.del(index)
|
||||||
|
#echo repr results
|
||||||
|
index = awaitAny(tasks)
|
||||||
|
|
||||||
|
doAssert results.len == 5
|
||||||
|
doAssert 1000 in results
|
||||||
|
doAssert 2000 in results
|
||||||
|
doAssert 3000 in results
|
||||||
|
doAssert 4000 in results
|
||||||
|
doAssert 5000 in results
|
||||||
|
sync()
|
||||||
|
echo "true"
|
||||||
Loading…
Add table
Add a link
Reference in a new issue