Refactored version of execProcesses with test.
This commit is contained in:
parent
2a7cfe4043
commit
c6c0d28a4f
2 changed files with 115 additions and 54 deletions
|
|
@ -231,55 +231,81 @@ proc execProcesses*(cmds: openArray[string],
|
||||||
## executes the commands `cmds` in parallel. Creates `n` processes
|
## executes the commands `cmds` in parallel. Creates `n` processes
|
||||||
## that execute in parallel. The highest return value of all processes
|
## that execute in parallel. The highest return value of all processes
|
||||||
## is returned. Runs `beforeRunEvent` before running each command.
|
## is returned. Runs `beforeRunEvent` before running each command.
|
||||||
when false:
|
|
||||||
# poParentStreams causes problems on Posix, so we simply disable it:
|
|
||||||
var options = options - {poParentStreams}
|
|
||||||
|
|
||||||
assert n > 0
|
assert n > 0
|
||||||
if n > 1:
|
if n > 1:
|
||||||
var q: seq[Process]
|
var i = 0
|
||||||
newSeq(q, n)
|
var q = newSeq[Process](n)
|
||||||
var m = min(n, cmds.len)
|
var m = min(n, cmds.len)
|
||||||
for i in 0..m-1:
|
|
||||||
|
when defined(windows):
|
||||||
|
var w: WOHandleArray
|
||||||
|
var wcount = m
|
||||||
|
for c in 0..MAXIMUM_WAIT_OBJECTS - 1:
|
||||||
|
w[c] = 0
|
||||||
|
|
||||||
|
while i < m:
|
||||||
if beforeRunEvent != nil:
|
if beforeRunEvent != nil:
|
||||||
beforeRunEvent(i)
|
beforeRunEvent(i)
|
||||||
q[i] = startProcess(cmds[i], options=options + {poEvalCommand})
|
q[i] = startProcess(cmds[i], options = options + {poEvalCommand})
|
||||||
when defined(noBusyWaiting):
|
when defined(windows):
|
||||||
var r = 0
|
w[i] = q[i].fProcessHandle
|
||||||
for i in m..high(cmds):
|
inc(i)
|
||||||
when defined(debugExecProcesses):
|
|
||||||
var err = ""
|
var ecount = len(cmds)
|
||||||
var outp = outputStream(q[r])
|
while ecount > 0:
|
||||||
while running(q[r]) or not atEnd(outp):
|
when defined(windows):
|
||||||
err.add(outp.readLine())
|
# waiting for all children, get result if any child exits
|
||||||
err.add("\n")
|
var ret = waitForMultipleObjects(int32(wcount), addr(w), 0'i32,
|
||||||
echo(err)
|
INFINITE)
|
||||||
result = max(waitForExit(q[r]), result)
|
if ret == WAIT_TIMEOUT:
|
||||||
if afterRunEvent != nil: afterRunEvent(r, q[r])
|
# must not be happen
|
||||||
if q[r] != nil: close(q[r])
|
discard
|
||||||
if beforeRunEvent != nil:
|
elif ret == WAIT_FAILED:
|
||||||
beforeRunEvent(i)
|
raiseOSError(osLastError())
|
||||||
q[r] = startProcess(cmds[i], options=options + {poEvalCommand})
|
else:
|
||||||
r = (r + 1) mod n
|
var status : cint = 1
|
||||||
else:
|
# waiting for all children, get result if any child exits
|
||||||
var i = m
|
let res = waitpid(-1, status, 0)
|
||||||
while i <= high(cmds):
|
if res > 0:
|
||||||
sleep(50)
|
for r in 0..m-1:
|
||||||
for r in 0..n-1:
|
if not isNil(q[r]) and q[r].id == res:
|
||||||
|
# we updating `exitStatus` manually, so `running()` can work.
|
||||||
|
if WIFEXITED(status) or WIFSIGNALED(status):
|
||||||
|
q[r].exitStatus = status
|
||||||
|
break
|
||||||
|
else:
|
||||||
|
let err = osLastError()
|
||||||
|
if err == OSErrorCode(ECHILD):
|
||||||
|
# some child exits, we need to check our childs exit codes
|
||||||
|
discard
|
||||||
|
elif err == OSErrorCode(EINTR):
|
||||||
|
# signal interrupted our syscall, lets repeat it
|
||||||
|
continue
|
||||||
|
else:
|
||||||
|
# all other errors are exceptions
|
||||||
|
raiseOSError(err)
|
||||||
|
|
||||||
|
for r in 0..m-1:
|
||||||
|
if not isNil(q[r]):
|
||||||
if not running(q[r]):
|
if not running(q[r]):
|
||||||
#echo(outputStream(q[r]).readLine())
|
result = max(result, q[r].peekExitCode())
|
||||||
result = max(waitForExit(q[r]), result)
|
|
||||||
if afterRunEvent != nil: afterRunEvent(r, q[r])
|
if afterRunEvent != nil: afterRunEvent(r, q[r])
|
||||||
if q[r] != nil: close(q[r])
|
close(q[r])
|
||||||
if beforeRunEvent != nil:
|
if i < len(cmds):
|
||||||
beforeRunEvent(i)
|
if beforeRunEvent != nil: beforeRunEvent(i)
|
||||||
q[r] = startProcess(cmds[i], options=options + {poEvalCommand})
|
q[r] = startProcess(cmds[i],
|
||||||
inc(i)
|
options = options + {poEvalCommand})
|
||||||
if i > high(cmds): break
|
when defined(windows):
|
||||||
for j in 0..m-1:
|
w[r] = q[r].fProcessHandle
|
||||||
result = max(waitForExit(q[j]), result)
|
inc(i)
|
||||||
if afterRunEvent != nil: afterRunEvent(j, q[j])
|
else:
|
||||||
if q[j] != nil: close(q[j])
|
q[r] = nil
|
||||||
|
when defined(windows):
|
||||||
|
for c in r..MAXIMUM_WAIT_OBJECTS - 2:
|
||||||
|
w[c] = w[c + 1]
|
||||||
|
dec(wcount)
|
||||||
|
dec(ecount)
|
||||||
else:
|
else:
|
||||||
for i in 0..high(cmds):
|
for i in 0..high(cmds):
|
||||||
if beforeRunEvent != nil:
|
if beforeRunEvent != nil:
|
||||||
|
|
@ -939,19 +965,22 @@ elif not defined(useNimRtl):
|
||||||
if kill(p.id, SIGCONT) != 0'i32: raiseOsError(osLastError())
|
if kill(p.id, SIGCONT) != 0'i32: raiseOsError(osLastError())
|
||||||
|
|
||||||
proc running(p: Process): bool =
|
proc running(p: Process): bool =
|
||||||
var ret : int
|
if p.exitStatus != -3:
|
||||||
var status : cint = 1
|
|
||||||
ret = waitpid(p.id, status, WNOHANG)
|
|
||||||
if ret == int(p.id):
|
|
||||||
if isExitStatus(status):
|
|
||||||
p.exitStatus = status
|
|
||||||
return false
|
|
||||||
else:
|
|
||||||
return true
|
|
||||||
elif ret == 0:
|
|
||||||
return true # Can't establish status. Assume running.
|
|
||||||
else:
|
|
||||||
return false
|
return false
|
||||||
|
else:
|
||||||
|
var ret : int
|
||||||
|
var status : cint = 1
|
||||||
|
ret = waitpid(p.id, status, WNOHANG)
|
||||||
|
if ret == int(p.id):
|
||||||
|
if isExitStatus(status):
|
||||||
|
p.exitStatus = status
|
||||||
|
return false
|
||||||
|
else:
|
||||||
|
return true
|
||||||
|
elif ret == 0:
|
||||||
|
return true # Can't establish status. Assume running.
|
||||||
|
else:
|
||||||
|
raiseOSError(osLastError())
|
||||||
|
|
||||||
proc terminate(p: Process) =
|
proc terminate(p: Process) =
|
||||||
if kill(p.id, SIGTERM) != 0'i32:
|
if kill(p.id, SIGTERM) != 0'i32:
|
||||||
|
|
|
||||||
32
tests/osproc/texecps.nim
Normal file
32
tests/osproc/texecps.nim
Normal file
|
|
@ -0,0 +1,32 @@
|
||||||
|
discard """
|
||||||
|
file: "texecps.nim"
|
||||||
|
output: ""
|
||||||
|
"""
|
||||||
|
|
||||||
|
import osproc, streams, strutils, os
|
||||||
|
|
||||||
|
const NumberOfProcesses = 13
|
||||||
|
|
||||||
|
var gResults {.threadvar.}: seq[string]
|
||||||
|
|
||||||
|
proc execCb(idx: int, p: Process) =
|
||||||
|
let exitCode = p.peekExitCode
|
||||||
|
if exitCode < len(gResults):
|
||||||
|
gResults[exitCode] = p.outputStream.readAll.strip
|
||||||
|
|
||||||
|
when isMainModule:
|
||||||
|
|
||||||
|
if paramCount() == 0:
|
||||||
|
gResults = newSeq[string](NumberOfProcesses)
|
||||||
|
var checks = newSeq[string](NumberOfProcesses)
|
||||||
|
var commands = newSeq[string](NumberOfProcesses)
|
||||||
|
for i in 0..len(commands) - 1:
|
||||||
|
commands[i] = getAppFileName() & " " & $i
|
||||||
|
checks[i] = $i
|
||||||
|
let cres = execProcesses(commands, options = {poStdErrToStdOut},
|
||||||
|
afterRunEvent = execCb)
|
||||||
|
doAssert(cres == len(commands) - 1)
|
||||||
|
doAssert(gResults == checks)
|
||||||
|
else:
|
||||||
|
echo paramStr(1)
|
||||||
|
programResult = parseInt(paramStr(1))
|
||||||
Loading…
Add table
Add a link
Reference in a new issue