Make sure an exception in any iterator gets raised in the main thread
Signed-off-by: Aanand Prasad <aanand.prasad@gmail.com> Conflicts: compose/cli/multiplexer.py
This commit is contained in:
parent
a9942b512a
commit
80d90a745a
2 changed files with 28 additions and 5 deletions
|
|
@ -26,7 +26,11 @@ class Multiplexer(object):
|
||||||
|
|
||||||
while self._num_running > 0:
|
while self._num_running > 0:
|
||||||
try:
|
try:
|
||||||
item = self.queue.get(timeout=0.1)
|
item, exception = self.queue.get(timeout=0.1)
|
||||||
|
|
||||||
|
if exception:
|
||||||
|
raise exception
|
||||||
|
|
||||||
if item is STOP:
|
if item is STOP:
|
||||||
self._num_running -= 1
|
self._num_running -= 1
|
||||||
else:
|
else:
|
||||||
|
|
@ -42,7 +46,9 @@ class Multiplexer(object):
|
||||||
|
|
||||||
|
|
||||||
def _enqueue_output(iterator, queue):
|
def _enqueue_output(iterator, queue):
|
||||||
|
try:
|
||||||
for item in iterator:
|
for item in iterator:
|
||||||
queue.put(item)
|
queue.put((item, None))
|
||||||
|
queue.put((STOP, None))
|
||||||
queue.put(STOP)
|
except Exception as e:
|
||||||
|
queue.put((None, e))
|
||||||
|
|
|
||||||
|
|
@ -26,3 +26,20 @@ class MultiplexerTest(unittest.TestCase):
|
||||||
[0, 1, 2, 3, 4, 5],
|
[0, 1, 2, 3, 4, 5],
|
||||||
sorted(list(mux.loop())),
|
sorted(list(mux.loop())),
|
||||||
)
|
)
|
||||||
|
|
||||||
|
def test_exception(self):
|
||||||
|
class Problem(Exception):
|
||||||
|
pass
|
||||||
|
|
||||||
|
def problematic_iterator():
|
||||||
|
yield 0
|
||||||
|
yield 2
|
||||||
|
raise Problem(":(")
|
||||||
|
|
||||||
|
mux = Multiplexer([
|
||||||
|
problematic_iterator(),
|
||||||
|
(x for x in [1, 3, 5]),
|
||||||
|
])
|
||||||
|
|
||||||
|
with self.assertRaises(Problem):
|
||||||
|
list(mux.loop())
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue