generated from 3amigos-dev/3amigos-py
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
- Loading branch information
Showing
3 changed files
with
168 additions
and
13 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,94 @@ | ||
""" | ||
Multithread processing to maximize time value of user input | ||
""" | ||
|
||
import concurrent.futures | ||
import logging | ||
|
||
|
||
class PoolManager: | ||
""" | ||
Used to add tasks that must be json serializable to pass to threads for | ||
execution or if requiring input saved to the user input priority heap. | ||
""" | ||
|
||
def __init__(self, handlers, max_workers): | ||
self._handlers = handlers | ||
self._executor = concurrent.futures.ThreadPoolExecutor(max_workers=max_workers) | ||
self._draining = False | ||
self._saved = [] | ||
|
||
def add(self, taskjson): | ||
""" | ||
add a task to the executor | ||
""" | ||
if self._draining: | ||
raise Exception("No new tasks when draining.") | ||
self._executor.submit(self.run_task, taskjson=taskjson) | ||
|
||
def run_task(self, taskjson): | ||
""" | ||
Called by a thread in the pool to run the task | ||
""" | ||
try: | ||
if self._draining: | ||
self._saved.append(taskjson) | ||
return | ||
handler = self.load_handler(taskjson) | ||
handler() | ||
except Exception: # pylint: disable=broad-except | ||
logging.exception("Unhandled error") | ||
|
||
def load_handler(self, taskjson): | ||
""" | ||
Lookup the handlers to return a task | ||
""" | ||
factory = self._handlers[taskjson["name"]] | ||
return factory(taskjson) | ||
|
||
def stop(self): | ||
""" | ||
Wait for current tasks to complete | ||
""" | ||
self._executor.shutdown() | ||
|
||
def __enter__(self): | ||
""" | ||
Implement python with interface | ||
""" | ||
return self | ||
|
||
def __exit__(self, type, value, traceback): # pylint: disable=redefined-builtin | ||
""" | ||
Implement python with interface | ||
""" | ||
self.stop() | ||
|
||
def drain(self): | ||
""" | ||
Signal to stop executing new tasks | ||
""" | ||
self._draining = True | ||
|
||
def save(self): | ||
""" | ||
Shutdown and collect and saved results | ||
""" | ||
self.drain() | ||
self.stop() | ||
return self._saved | ||
|
||
|
||
def get_pool(handlers, max_workers=5): | ||
""" | ||
Obtain a threadpool | ||
""" | ||
return PoolManager(handlers, max_workers) | ||
|
||
|
||
def main(handlers): | ||
""" | ||
Main entry point to run pool manager | ||
""" | ||
with get_pool(handlers) as pool: | ||
return pool.save() |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,64 @@ | ||
""" | ||
Test cases to ensure tasks are picked up and executed concurrently whilst | ||
serializing user input | ||
""" | ||
|
||
import threading | ||
|
||
from meticulous._threadpool import get_pool | ||
|
||
|
||
def test_add_async(): | ||
""" | ||
Check adding async task correctly runs | ||
""" | ||
# Setup | ||
result = [] | ||
|
||
def run(): | ||
result.append(True) | ||
return result | ||
|
||
def load_run(_): | ||
return run | ||
|
||
pool = get_pool({"run": load_run}) | ||
# Exercise | ||
pool.add({"name": "run"}) | ||
# Verify | ||
pool.stop() | ||
assert result[0] # noqa=S101 # nosec | ||
|
||
|
||
def test_shutdown(): | ||
""" | ||
Check saving async task beyond number of works suspends correctly | ||
""" | ||
# Setup | ||
cond = threading.Condition() | ||
running = [0] | ||
|
||
def run(): | ||
with cond: | ||
running[0] += 1 | ||
cond.wait(60) | ||
|
||
def load_run(_): | ||
return run | ||
|
||
pool = get_pool({"run": load_run}, max_workers=2) | ||
taskjson = {"name": "run"} | ||
for _ in range(10): | ||
pool.add(taskjson) | ||
with cond: | ||
while running[0] < 2: | ||
cond.wait() | ||
pool.drain() | ||
with cond: | ||
cond.notify() | ||
cond.notify() | ||
pool.stop() | ||
# Exercise | ||
result = pool.save() | ||
# Verify | ||
assert result == ([taskjson] * 8) # noqa=S101 # nosec |