2016-12-08 20:48:59 +00:00
|
|
|
# -*- coding: utf-8 -*-
|
|
|
|
#
|
2017-07-26 15:39:07 +00:00
|
|
|
# 2016-2017 Darko Poljak (darko.poljak at gmail.com)
|
2016-12-08 20:48:59 +00:00
|
|
|
#
|
|
|
|
# This file is part of cdist.
|
|
|
|
#
|
|
|
|
# cdist is free software: you can redistribute it and/or modify
|
|
|
|
# it under the terms of the GNU General Public License as published by
|
|
|
|
# the Free Software Foundation, either version 3 of the License, or
|
|
|
|
# (at your option) any later version.
|
|
|
|
#
|
|
|
|
# cdist is distributed in the hope that it will be useful,
|
|
|
|
# but WITHOUT ANY WARRANTY; without even the implied warranty of
|
|
|
|
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
|
|
|
# GNU General Public License for more details.
|
|
|
|
#
|
|
|
|
# You should have received a copy of the GNU General Public License
|
|
|
|
# along with cdist. If not, see <http://www.gnu.org/licenses/>.
|
|
|
|
#
|
|
|
|
#
|
|
|
|
|
|
|
|
|
|
|
|
import multiprocessing
|
2017-07-26 10:01:19 +00:00
|
|
|
import concurrent.futures as cf
|
2016-12-08 20:48:59 +00:00
|
|
|
import itertools
|
2017-07-26 10:01:19 +00:00
|
|
|
import os
|
|
|
|
import signal
|
|
|
|
import logging
|
|
|
|
|
|
|
|
|
|
|
|
log = logging.getLogger("cdist-mputil")
|
2016-12-08 20:48:59 +00:00
|
|
|
|
|
|
|
|
2017-07-26 15:39:07 +00:00
|
|
|
def mp_sig_handler(signum, frame):
|
|
|
|
log.trace("signal %s, SIGKILL whole process group", signum)
|
|
|
|
os.killpg(os.getpgrp(), signal.SIGKILL)
|
|
|
|
|
|
|
|
|
2016-12-08 20:48:59 +00:00
|
|
|
def mp_pool_run(func, args=None, kwds=None, jobs=multiprocessing.cpu_count()):
|
2017-07-26 10:01:19 +00:00
|
|
|
"""Run func using concurrent.futures.ProcessPoolExecutor with jobs jobs
|
|
|
|
and supplied iterables of args and kwds with one entry for each
|
|
|
|
parallel func instance.
|
|
|
|
Return list of results.
|
2016-12-08 20:48:59 +00:00
|
|
|
"""
|
|
|
|
if args and kwds:
|
2016-12-11 20:17:22 +00:00
|
|
|
fargs = zip(args, kwds)
|
2016-12-08 20:48:59 +00:00
|
|
|
elif args:
|
|
|
|
fargs = zip(args, itertools.repeat({}))
|
|
|
|
elif kwds:
|
|
|
|
fargs = zip(itertools.repeat(()), kwds)
|
|
|
|
else:
|
|
|
|
return [func(), ]
|
|
|
|
|
2017-07-26 10:01:19 +00:00
|
|
|
retval = []
|
|
|
|
with cf.ProcessPoolExecutor(jobs) as executor:
|
|
|
|
try:
|
|
|
|
results = [
|
|
|
|
executor.submit(func, *a, **k) for a, k in fargs
|
|
|
|
]
|
|
|
|
for f in cf.as_completed(results):
|
|
|
|
retval.append(f.result())
|
|
|
|
return retval
|
|
|
|
except KeyboardInterrupt:
|
2017-07-26 15:39:07 +00:00
|
|
|
mp_sig_handler(signal.SIGINT, None)
|
2017-07-26 10:01:19 +00:00
|
|
|
raise
|