"""Logic for executing other logic layer methods in parallel.""" import concurrent.futures from oto import response as oto_response from sound_recordings.api import executor REQUEST_ERROR = {"error": "unexpected error"} def _format_response(response): """Handle formatting of responses from requests.""" if not hasattr(response, "status"): return response if response.status != 200: return REQUEST_ERROR return response.message def parallel(requests): """Execute requests in parallel. Provided a dict of requests, execute them with the provided id and kwargs. Return the results of these requests in a dict keyed by keys in the requests parameter. """ futures = { executor.submit(request["func"], *request["args"]): key for (key, request) in requests.items() } concurrent.futures.wait(futures) output = dict(map(lambda f: (futures[f], f.result()), futures)) # Format the responses for errors, return as oto response return oto_response.Response( dict(map(lambda kv: (kv[0], _format_response(kv[1])), output.items())) )